syncMain.c 121.4 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"
M
Minghao Li 已提交
26
#include "syncRaftCfg.h"
M
Minghao Li 已提交
27
#include "syncRaftLog.h"
M
Minghao Li 已提交
28
#include "syncRaftStore.h"
M
Minghao Li 已提交
29
#include "syncReplication.h"
M
Minghao Li 已提交
30 31
#include "syncRequestVote.h"
#include "syncRequestVoteReply.h"
M
Minghao Li 已提交
32
#include "syncRespMgr.h"
M
Minghao Li 已提交
33
#include "syncSnapshot.h"
M
Minghao Li 已提交
34
#include "syncTimeout.h"
M
Minghao Li 已提交
35
#include "syncUtil.h"
M
Minghao Li 已提交
36
#include "syncVoteMgr.h"
M
Minghao Li 已提交
37

M
Minghao Li 已提交
38 39 40 41 42
static void    syncNodeEqPingTimer(void* param, void* tmrId);
static void    syncNodeEqElectTimer(void* param, void* tmrId);
static void    syncNodeEqHeartbeatTimer(void* param, void* tmrId);
static int32_t syncNodeEqNoop(SSyncNode* ths);
static int32_t syncNodeAppendNoop(SSyncNode* ths);
43
static void    syncNodeEqPeerHeartbeatTimer(void* param, void* tmrId);
S
Shengliang Guan 已提交
44
static bool    syncIsConfigChanged(const SSyncCfg* pOldCfg, const SSyncCfg* pNewCfg);
M
Minghao Li 已提交
45

46
int64_t syncOpen(SSyncInfo* pSyncInfo) {
S
Shengliang Guan 已提交
47 48
  SSyncNode* pSyncNode = syncNodeOpen(pSyncInfo);
  if (pSyncNode == NULL) {
S
Shengliang Guan 已提交
49
    sError("vgId:%d, failed to open sync node", pSyncInfo->vgId);
50 51
    return -1;
  }
M
Minghao Li 已提交
52

S
Shengliang Guan 已提交
53 54 55
  pSyncNode->rid = syncNodeAdd(pSyncNode);
  if (pSyncNode->rid < 0) {
    syncNodeClose(pSyncNode);
M
Minghao Li 已提交
56 57 58
    return -1;
  }

S
Shengliang Guan 已提交
59 60 61 62 63 64
  pSyncNode->pingBaseLine = pSyncInfo->pingMs;
  pSyncNode->pingTimerMS = pSyncInfo->pingMs;
  pSyncNode->electBaseLine = pSyncInfo->electMs;
  pSyncNode->hbBaseLine = pSyncInfo->heartbeatMs;
  pSyncNode->heartbeatTimerMS = pSyncInfo->heartbeatMs;
  pSyncNode->msgcb = pSyncInfo->msgcb;
S
Shengliang Guan 已提交
65
  return pSyncNode->rid;
M
Minghao Li 已提交
66
}
M
Minghao Li 已提交
67

M
Minghao Li 已提交
68
void syncStart(int64_t rid) {
S
Shengliang Guan 已提交
69 70 71 72
  SSyncNode* pSyncNode = syncNodeAcquire(rid);
  if (pSyncNode != NULL) {
    syncNodeStart(pSyncNode);
    syncNodeRelease(pSyncNode);
M
Minghao Li 已提交
73 74 75
  }
}

M
Minghao Li 已提交
76
void syncStop(int64_t rid) {
S
Shengliang Guan 已提交
77
  SSyncNode* pSyncNode = syncNodeAcquire(rid);
S
Shengliang Guan 已提交
78
  if (pSyncNode != NULL) {
S
Shengliang Guan 已提交
79
    syncNodeRelease(pSyncNode);
S
Shengliang Guan 已提交
80
    syncNodeRemove(rid);
M
Minghao Li 已提交
81
  }
S
Shengliang Guan 已提交
82
}
M
Minghao Li 已提交
83

S
Shengliang Guan 已提交
84 85 86
static bool syncNodeCheckNewConfig(SSyncNode* pSyncNode, const SSyncCfg* pCfg) {
  if (!syncNodeInConfig(pSyncNode, pCfg)) return false;
  return abs(pCfg->replicaNum - pSyncNode->replicaNum) <= 1;
M
Minghao Li 已提交
87 88
}

S
Shengliang Guan 已提交
89
int32_t syncReconfig(int64_t rid, SSyncCfg* pNewCfg) {
S
Shengliang Guan 已提交
90
  SSyncNode* pSyncNode = syncNodeAcquire(rid);
S
Shengliang Guan 已提交
91
  if (pSyncNode == NULL) return -1;
M
Minghao Li 已提交
92

M
Minghao Li 已提交
93
  if (!syncNodeCheckNewConfig(pSyncNode, pNewCfg)) {
S
Shengliang Guan 已提交
94
    syncNodeRelease(pSyncNode);
M
Minghao Li 已提交
95
    terrno = TSDB_CODE_SYN_NEW_CONFIG_ERROR;
S
Shengliang Guan 已提交
96
    sError("vgId:%d, failed to reconfig since invalid new config", pSyncNode->vgId);
M
Minghao Li 已提交
97
    return -1;
M
Minghao Li 已提交
98
  }
99

S
Shengliang Guan 已提交
100 101
  syncNodeUpdateNewConfigIndex(pSyncNode, pNewCfg);
  syncNodeDoConfigChange(pSyncNode, pNewCfg, SYNC_INDEX_INVALID);
S
Shengliang Guan 已提交
102

M
Minghao Li 已提交
103 104 105 106
  if (pSyncNode->state == TAOS_SYNC_STATE_LEADER) {
    syncNodeStopHeartbeatTimer(pSyncNode);

    for (int32_t i = 0; i < TSDB_MAX_REPLICA; ++i) {
S
Shengliang Guan 已提交
107
      syncHbTimerInit(pSyncNode, &pSyncNode->peerHeartbeatTimerArr[i], pSyncNode->replicasId[i]);
M
Minghao Li 已提交
108 109 110 111 112
    }

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

S
Shengliang Guan 已提交
114
  syncNodeRelease(pSyncNode);
S
Shengliang Guan 已提交
115
  return 0;
M
Minghao Li 已提交
116
}
M
Minghao Li 已提交
117

S
Shengliang Guan 已提交
118 119 120 121
int32_t syncProcessMsg(int64_t rid, SRpcMsg* pMsg) {
  int32_t code = -1;
  if (!syncIsInit()) return code;

S
Shengliang Guan 已提交
122
  SSyncNode* pSyncNode = syncNodeAcquire(rid);
S
Shengliang Guan 已提交
123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179
  if (pSyncNode == NULL) return code;

  if (pMsg->msgType == TDMT_SYNC_HEARTBEAT) {
    SyncHeartbeat* pSyncMsg = syncHeartbeatFromRpcMsg2(pMsg);
    code = syncNodeOnHeartbeat(pSyncNode, pSyncMsg);
    syncHeartbeatDestroy(pSyncMsg);
  } else if (pMsg->msgType == TDMT_SYNC_HEARTBEAT_REPLY) {
    SyncHeartbeatReply* pSyncMsg = syncHeartbeatReplyFromRpcMsg2(pMsg);
    code = syncNodeOnHeartbeatReply(pSyncNode, pSyncMsg);
    syncHeartbeatReplyDestroy(pSyncMsg);
  } else if (pMsg->msgType == TDMT_SYNC_TIMEOUT) {
    SyncTimeout* pSyncMsg = syncTimeoutFromRpcMsg2(pMsg);
    code = syncNodeOnTimer(pSyncNode, pSyncMsg);
    syncTimeoutDestroy(pSyncMsg);
  } else if (pMsg->msgType == TDMT_SYNC_PING) {
    SyncPing* pSyncMsg = syncPingFromRpcMsg2(pMsg);
    code = syncNodeOnPing(pSyncNode, pSyncMsg);
    syncPingDestroy(pSyncMsg);
  } else if (pMsg->msgType == TDMT_SYNC_PING_REPLY) {
    SyncPingReply* pSyncMsg = syncPingReplyFromRpcMsg2(pMsg);
    code = syncNodeOnPingReply(pSyncNode, pSyncMsg);
    syncPingReplyDestroy(pSyncMsg);
  } else if (pMsg->msgType == TDMT_SYNC_CLIENT_REQUEST) {
    SyncClientRequest* pSyncMsg = syncClientRequestFromRpcMsg2(pMsg);
    code = syncNodeOnClientRequest(pSyncNode, pSyncMsg, NULL);
    syncClientRequestDestroy(pSyncMsg);
  } else if (pMsg->msgType == TDMT_SYNC_REQUEST_VOTE) {
    SyncRequestVote* pSyncMsg = syncRequestVoteFromRpcMsg2(pMsg);
    code = syncNodeOnRequestVote(pSyncNode, pSyncMsg);
    syncRequestVoteDestroy(pSyncMsg);
  } else if (pMsg->msgType == TDMT_SYNC_REQUEST_VOTE_REPLY) {
    SyncRequestVoteReply* pSyncMsg = syncRequestVoteReplyFromRpcMsg2(pMsg);
    code = syncNodeOnRequestVoteReply(pSyncNode, pSyncMsg);
    syncRequestVoteReplyDestroy(pSyncMsg);
  } else if (pMsg->msgType == TDMT_SYNC_APPEND_ENTRIES) {
    SyncAppendEntries* pSyncMsg = syncAppendEntriesFromRpcMsg2(pMsg);
    code = syncNodeOnAppendEntries(pSyncNode, pSyncMsg);
    syncAppendEntriesDestroy(pSyncMsg);
  } else if (pMsg->msgType == TDMT_SYNC_APPEND_ENTRIES_REPLY) {
    SyncAppendEntriesReply* pSyncMsg = syncAppendEntriesReplyFromRpcMsg2(pMsg);
    code = syncNodeOnAppendEntriesReply(pSyncNode, pSyncMsg);
    syncAppendEntriesReplyDestroy(pSyncMsg);
  } else if (pMsg->msgType == TDMT_SYNC_SNAPSHOT_SEND) {
    SyncSnapshotSend* pSyncMsg = syncSnapshotSendFromRpcMsg2(pMsg);
    code = syncNodeOnSnapshot(pSyncNode, pSyncMsg);
    syncSnapshotSendDestroy(pSyncMsg);
  } else if (pMsg->msgType == TDMT_SYNC_SNAPSHOT_RSP) {
    SyncSnapshotRsp* pSyncMsg = syncSnapshotRspFromRpcMsg2(pMsg);
    code = syncNodeOnSnapshotReply(pSyncNode, pSyncMsg);
    syncSnapshotRspDestroy(pSyncMsg);
  } else if (pMsg->msgType == TDMT_SYNC_LOCAL_CMD) {
    SyncLocalCmd* pSyncMsg = syncLocalCmdFromRpcMsg2(pMsg);
    code = syncNodeOnLocalCmd(pSyncNode, pSyncMsg);
    syncLocalCmdDestroy(pSyncMsg);
  } else {
    sError("vgId:%d, failed to process msg:%p since invalid type:%s", pSyncNode->vgId, pMsg, TMSG_INFO(pMsg->msgType));
    code = -1;
M
Minghao Li 已提交
180 181
  }

S
Shengliang Guan 已提交
182
  syncNodeRelease(pSyncNode);
S
Shengliang Guan 已提交
183
  return code;
184 185
}

S
Shengliang Guan 已提交
186
int32_t syncLeaderTransfer(int64_t rid) {
S
Shengliang Guan 已提交
187
  SSyncNode* pSyncNode = syncNodeAcquire(rid);
S
Shengliang Guan 已提交
188
  if (pSyncNode == NULL) return -1;
189

S
Shengliang Guan 已提交
190
  int32_t ret = syncNodeLeaderTransfer(pSyncNode);
S
Shengliang Guan 已提交
191
  syncNodeRelease(pSyncNode);
192 193 194
  return ret;
}

M
Minghao Li 已提交
195 196 197 198 199 200 201 202 203 204 205 206 207 208 209 210
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;
}

211 212 213 214 215 216 217 218 219 220 221 222 223 224
char* syncNodePeerState2Str(const SSyncNode* pSyncNode) {
  int32_t len = 128;
  int32_t useLen = 0;
  int32_t leftLen = len - useLen;
  char*   pStr = taosMemoryMalloc(len);
  memset(pStr, 0, len);

  char*   p = pStr;
  int32_t use = snprintf(p, leftLen, "{");
  useLen += use;
  leftLen -= use;

  for (int32_t i = 0; i < pSyncNode->replicaNum; ++i) {
    SPeerState* pState = syncNodeGetPeerState((SSyncNode*)pSyncNode, &(pSyncNode->replicasId[i]));
M
Minghao Li 已提交
225
    if (pState == NULL) {
226
      sError("vgId:%d, replica maybe dropped", pSyncNode->vgId);
M
Minghao Li 已提交
227 228
      break;
    }
229 230

    p = pStr + useLen;
S
Shengliang Guan 已提交
231
    use = snprintf(p, leftLen, "%d:%" PRId64 " ,%" PRId64, i, pState->lastSendIndex, pState->lastSendTime);
232 233 234 235 236 237 238 239 240 241 242 243 244 245
    useLen += use;
    leftLen -= use;
  }

  p = pStr + useLen;
  use = snprintf(p, leftLen, "}");
  useLen += use;
  leftLen -= use;

  // sTrace("vgId:%d, ------------------ syncNodePeerState2Str:%s", pSyncNode->vgId, pStr);

  return pStr;
}

246
int32_t syncBeginSnapshot(int64_t rid, int64_t lastApplyIndex) {
S
Shengliang Guan 已提交
247
  SSyncNode* pSyncNode = syncNodeAcquire(rid);
248 249 250 251 252 253 254
  if (pSyncNode == NULL) {
    terrno = TSDB_CODE_SYN_INTERNAL_ERROR;
    return -1;
  }
  ASSERT(rid == pSyncNode->rid);
  int32_t code = 0;

M
Minghao Li 已提交
255
  if (syncNodeIsMnode(pSyncNode)) {
M
Minghao Li 已提交
256 257 258
    // mnode
    int64_t logRetention = SYNC_MNODE_LOG_RETENTION;

M
Minghao Li 已提交
259 260 261
    SyncIndex beginIndex = pSyncNode->pLogStore->syncLogBeginIndex(pSyncNode->pLogStore);
    SyncIndex endIndex = pSyncNode->pLogStore->syncLogEndIndex(pSyncNode->pLogStore);
    int64_t   logNum = endIndex - beginIndex;
M
Minghao Li 已提交
262 263 264
    bool      isEmpty = pSyncNode->pLogStore->syncLogIsEmpty(pSyncNode->pLogStore);

    if (isEmpty || (!isEmpty && logNum < logRetention)) {
M
Minghao Li 已提交
265
      char logBuf[256];
S
Shengliang Guan 已提交
266 267 268
      snprintf(logBuf, sizeof(logBuf),
               "new-snapshot-index:%" PRId64 ", log-num:%" PRId64 ", empty:%d, do not delete wal", lastApplyIndex,
               logNum, isEmpty);
M
Minghao Li 已提交
269 270
      syncNodeEventLog(pSyncNode, logBuf);

S
Shengliang Guan 已提交
271
      syncNodeRelease(pSyncNode);
M
Minghao Li 已提交
272 273 274
      return 0;
    }

M
Minghao Li 已提交
275 276 277 278 279 280 281 282 283 284 285 286 287 288 289 290 291 292 293
    goto _DEL_WAL;

  } else {
    // vnode
    if (pSyncNode->replicaNum > 1) {
      // multi replicas

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

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

S
Shengliang Guan 已提交
300
            syncNodeRelease(pSyncNode);
M
Minghao Li 已提交
301 302 303 304 305 306 307 308
            return 0;
          }
        }

      } else if (pSyncNode->state == TAOS_SYNC_STATE_FOLLOWER) {
        if (lastApplyIndex > pSyncNode->minMatchIndex) {
          char logBuf[256];
          snprintf(logBuf, sizeof(logBuf),
S
Shengliang Guan 已提交
309 310
                   "new-snapshot-index:%" PRId64 " is greater than min-match-index:%" PRId64 ", do not delete wal",
                   lastApplyIndex, pSyncNode->minMatchIndex);
M
Minghao Li 已提交
311 312
          syncNodeEventLog(pSyncNode, logBuf);

S
Shengliang Guan 已提交
313
          syncNodeRelease(pSyncNode);
M
Minghao Li 已提交
314 315 316 317
          return 0;
        }

      } else if (pSyncNode->state == TAOS_SYNC_STATE_CANDIDATE) {
318
        char logBuf[256];
S
Shengliang Guan 已提交
319
        snprintf(logBuf, sizeof(logBuf), "new-snapshot-index:%" PRId64 " candidate, do not delete wal", lastApplyIndex);
320 321
        syncNodeEventLog(pSyncNode, logBuf);

S
Shengliang Guan 已提交
322
        syncNodeRelease(pSyncNode);
M
Minghao Li 已提交
323 324 325 326
        return 0;

      } else {
        char logBuf[256];
S
Shengliang Guan 已提交
327 328
        snprintf(logBuf, sizeof(logBuf), "new-snapshot-index:%" PRId64 " unknown state, do not delete wal",
                 lastApplyIndex);
M
Minghao Li 已提交
329 330
        syncNodeEventLog(pSyncNode, logBuf);

S
Shengliang Guan 已提交
331
        syncNodeRelease(pSyncNode);
M
Minghao Li 已提交
332 333 334 335 336 337 338 339 340
        return 0;
      }

      goto _DEL_WAL;

    } else {
      // one replica

      goto _DEL_WAL;
341 342 343
    }
  }

M
Minghao Li 已提交
344
_DEL_WAL:
345

M
Minghao Li 已提交
346
  do {
347 348 349 350
    SyncIndex snapshottingIndex = atomic_load_64(&pSyncNode->snapshottingIndex);

    if (snapshottingIndex == SYNC_INDEX_INVALID) {
      atomic_store_64(&pSyncNode->snapshottingIndex, lastApplyIndex);
M
Minghao Li 已提交
351
      pSyncNode->snapshottingTime = taosGetTimestampMs();
352

M
Minghao Li 已提交
353 354 355
      SSyncLogStoreData* pData = pSyncNode->pLogStore->data;
      code = walBeginSnapshot(pData->pWal, lastApplyIndex);
      if (code == 0) {
356
        char logBuf[256];
S
Shengliang Guan 已提交
357
        snprintf(logBuf, sizeof(logBuf), "wal snapshot begin, index:%" PRId64 ", last apply index:%" PRId64,
358 359 360
                 pSyncNode->snapshottingIndex, lastApplyIndex);
        syncNodeEventLog(pSyncNode, logBuf);

M
Minghao Li 已提交
361 362
      } else {
        char logBuf[256];
S
Shengliang Guan 已提交
363 364 365
        snprintf(logBuf, sizeof(logBuf),
                 "wal snapshot begin error since:%s, index:%" PRId64 ", last apply index:%" PRId64, terrstr(terrno),
                 pSyncNode->snapshottingIndex, lastApplyIndex);
M
Minghao Li 已提交
366 367 368 369
        syncNodeErrorLog(pSyncNode, logBuf);

        atomic_store_64(&pSyncNode->snapshottingIndex, SYNC_INDEX_INVALID);
      }
370 371

    } else {
372
      char logBuf[256];
S
Shengliang Guan 已提交
373 374 375
      snprintf(logBuf, sizeof(logBuf),
               "snapshotting for %" PRId64 ", do not delete wal for new-snapshot-index:%" PRId64, snapshottingIndex,
               lastApplyIndex);
376
      syncNodeEventLog(pSyncNode, logBuf);
377
    }
M
Minghao Li 已提交
378
  } while (0);
379

S
Shengliang Guan 已提交
380
  syncNodeRelease(pSyncNode);
381 382 383 384
  return code;
}

int32_t syncEndSnapshot(int64_t rid) {
S
Shengliang Guan 已提交
385
  SSyncNode* pSyncNode = syncNodeAcquire(rid);
386 387 388 389 390 391
  if (pSyncNode == NULL) {
    terrno = TSDB_CODE_SYN_INTERNAL_ERROR;
    return -1;
  }
  ASSERT(rid == pSyncNode->rid);

392 393 394 395
  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 已提交
396
    if (code != 0) {
397
      sError("vgId:%d, wal snapshot end error since:%s", pSyncNode->vgId, terrstr());
M
Minghao Li 已提交
398

S
Shengliang Guan 已提交
399
      syncNodeRelease(pSyncNode);
M
Minghao Li 已提交
400 401 402 403
      return -1;
    } else {
      do {
        char logBuf[256];
S
Shengliang Guan 已提交
404 405
        snprintf(logBuf, sizeof(logBuf), "wal snapshot end, index:%" PRId64,
                 atomic_load_64(&pSyncNode->snapshottingIndex));
M
Minghao Li 已提交
406 407
        syncNodeEventLog(pSyncNode, logBuf);
      } while (0);
408

M
Minghao Li 已提交
409 410
      atomic_store_64(&pSyncNode->snapshottingIndex, SYNC_INDEX_INVALID);
    }
411
  }
412

S
Shengliang Guan 已提交
413
  syncNodeRelease(pSyncNode);
414 415 416
  return code;
}

M
Minghao Li 已提交
417
int32_t syncStepDown(int64_t rid, SyncTerm newTerm) {
S
Shengliang Guan 已提交
418
  SSyncNode* pSyncNode = syncNodeAcquire(rid);
M
Minghao Li 已提交
419 420 421 422 423 424 425 426
  if (pSyncNode == NULL) {
    terrno = TSDB_CODE_SYN_INTERNAL_ERROR;
    return -1;
  }
  ASSERT(rid == pSyncNode->rid);

  syncNodeStepDown(pSyncNode, newTerm);

S
Shengliang Guan 已提交
427
  syncNodeRelease(pSyncNode);
M
Minghao Li 已提交
428 429 430
  return 0;
}

M
Minghao Li 已提交
431 432
int32_t syncNodeLeaderTransfer(SSyncNode* pSyncNode) {
  if (pSyncNode->peersNum == 0) {
433
    sDebug("only one replica, cannot leader transfer");
M
Minghao Li 已提交
434 435 436 437 438 439 440 441 442 443 444 445 446
    terrno = TSDB_CODE_SYN_ONE_REPLICA;
    return -1;
  }

  SNodeInfo newLeader = (pSyncNode->peersNodeInfo)[0];
  int32_t   ret = syncNodeLeaderTransferTo(pSyncNode, newLeader);
  return ret;
}

int32_t syncNodeLeaderTransferTo(SSyncNode* pSyncNode, SNodeInfo newLeader) {
  int32_t ret = 0;

  if (pSyncNode->replicaNum == 1) {
447
    sDebug("only one replica, cannot leader transfer");
M
Minghao Li 已提交
448 449 450 451
    terrno = TSDB_CODE_SYN_ONE_REPLICA;
    return -1;
  }

M
Minghao Li 已提交
452 453 454 455 456 457
  do {
    char logBuf[128];
    snprintf(logBuf, sizeof(logBuf), "begin leader transfer to %s:%u", newLeader.nodeFqdn, newLeader.nodePort);
    syncNodeEventLog(pSyncNode, logBuf);
  } while (0);

M
Minghao Li 已提交
458 459 460 461 462 463 464 465 466 467 468 469 470
  SyncLeaderTransfer* pMsg = syncLeaderTransferBuild(pSyncNode->vgId);
  pMsg->newLeaderId.addr = syncUtilAddr2U64(newLeader.nodeFqdn, newLeader.nodePort);
  pMsg->newLeaderId.vgId = pSyncNode->vgId;
  pMsg->newNodeInfo = newLeader;
  ASSERT(pMsg != NULL);
  SRpcMsg rpcMsg = {0};
  syncLeaderTransfer2RpcMsg(pMsg, &rpcMsg);
  syncLeaderTransferDestroy(pMsg);

  ret = syncNodePropose(pSyncNode, &rpcMsg, false);
  return ret;
}

471
bool syncCanLeaderTransfer(int64_t rid) {
S
Shengliang Guan 已提交
472
  SSyncNode* pSyncNode = syncNodeAcquire(rid);
473 474 475
  if (pSyncNode == NULL) {
    return false;
  }
M
Minghao Li 已提交
476
  ASSERT(rid == pSyncNode->rid);
477 478

  if (pSyncNode->replicaNum == 1) {
S
Shengliang Guan 已提交
479
    syncNodeRelease(pSyncNode);
480 481 482 483
    return false;
  }

  if (pSyncNode->state == TAOS_SYNC_STATE_FOLLOWER) {
S
Shengliang Guan 已提交
484
    syncNodeRelease(pSyncNode);
485 486 487 488 489 490 491 492 493 494 495 496 497 498
    return true;
  }

  bool matchOK = true;
  if (pSyncNode->state == TAOS_SYNC_STATE_CANDIDATE || pSyncNode->state == TAOS_SYNC_STATE_LEADER) {
    SyncIndex myCommitIndex = pSyncNode->commitIndex;
    for (int i = 0; i < pSyncNode->peersNum; ++i) {
      SyncIndex peerMatchIndex = syncIndexMgrGetIndex(pSyncNode->pMatchIndex, &(pSyncNode->peersId)[i]);
      if (peerMatchIndex < myCommitIndex) {
        matchOK = false;
      }
    }
  }

S
Shengliang Guan 已提交
499
  syncNodeRelease(pSyncNode);
500 501 502
  return matchOK;
}

503
int32_t syncForwardToPeer(int64_t rid, SRpcMsg* pMsg, bool isWeak) {
M
Minghao Li 已提交
504 505 506
  int32_t ret = syncPropose(rid, pMsg, isWeak);
  return ret;
}
M
Minghao Li 已提交
507

M
Minghao Li 已提交
508
ESyncState syncGetMyRole(int64_t rid) {
S
Shengliang Guan 已提交
509
  SSyncNode* pSyncNode = syncNodeAcquire(rid);
M
Minghao Li 已提交
510 511 512
  if (pSyncNode == NULL) {
    return TAOS_SYNC_STATE_ERROR;
  }
M
Minghao Li 已提交
513
  ASSERT(rid == pSyncNode->rid);
M
Minghao Li 已提交
514 515
  ESyncState state = pSyncNode->state;

S
Shengliang Guan 已提交
516
  syncNodeRelease(pSyncNode);
M
Minghao Li 已提交
517
  return state;
M
Minghao Li 已提交
518 519
}

M
Minghao Li 已提交
520
bool syncIsReady(int64_t rid) {
S
Shengliang Guan 已提交
521
  SSyncNode* pSyncNode = syncNodeAcquire(rid);
M
Minghao Li 已提交
522 523 524
  if (pSyncNode == NULL) {
    return false;
  }
M
Minghao Li 已提交
525
  ASSERT(rid == pSyncNode->rid);
M
Minghao Li 已提交
526
  bool b = (pSyncNode->state == TAOS_SYNC_STATE_LEADER) && pSyncNode->restoreFinish;
S
Shengliang Guan 已提交
527
  syncNodeRelease(pSyncNode);
528 529 530 531 532 533 534 535 536

  // if false, set error code
  if (false == b) {
    if (pSyncNode->state != TAOS_SYNC_STATE_LEADER) {
      terrno = TSDB_CODE_SYN_NOT_LEADER;
    } else {
      terrno = TSDB_CODE_APP_NOT_READY;
    }
  }
M
Minghao Li 已提交
537 538 539
  return b;
}

M
Minghao Li 已提交
540
bool syncIsRestoreFinish(int64_t rid) {
S
Shengliang Guan 已提交
541
  SSyncNode* pSyncNode = syncNodeAcquire(rid);
M
Minghao Li 已提交
542 543 544
  if (pSyncNode == NULL) {
    return false;
  }
M
Minghao Li 已提交
545
  ASSERT(rid == pSyncNode->rid);
M
Minghao Li 已提交
546 547
  bool b = pSyncNode->restoreFinish;

S
Shengliang Guan 已提交
548
  syncNodeRelease(pSyncNode);
M
Minghao Li 已提交
549 550 551
  return b;
}

552 553 554 555 556
int32_t syncGetSnapshotByIndex(int64_t rid, SyncIndex index, SSnapshot* pSnapshot) {
  if (index < SYNC_INDEX_BEGIN) {
    return -1;
  }

S
Shengliang Guan 已提交
557
  SSyncNode* pSyncNode = syncNodeAcquire(rid);
558 559 560 561 562 563 564 565 566 567 568
  if (pSyncNode == NULL) {
    return -1;
  }
  ASSERT(rid == pSyncNode->rid);

  SSyncRaftEntry* pEntry = NULL;
  int32_t         code = pSyncNode->pLogStore->syncLogGetEntry(pSyncNode->pLogStore, index, &pEntry);
  if (code != 0) {
    if (pEntry != NULL) {
      syncEntryDestory(pEntry);
    }
S
Shengliang Guan 已提交
569
    syncNodeRelease(pSyncNode);
570 571 572 573 574 575 576 577 578 579
    return -1;
  }
  ASSERT(pEntry != NULL);

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

  syncEntryDestory(pEntry);
S
Shengliang Guan 已提交
580
  syncNodeRelease(pSyncNode);
581 582 583
  return 0;
}

584
int32_t syncGetSnapshotMeta(int64_t rid, struct SSnapshotMeta* sMeta) {
S
Shengliang Guan 已提交
585
  SSyncNode* pSyncNode = syncNodeAcquire(rid);
586 587 588
  if (pSyncNode == NULL) {
    return -1;
  }
M
Minghao Li 已提交
589
  ASSERT(rid == pSyncNode->rid);
590 591
  sMeta->lastConfigIndex = pSyncNode->pRaftCfg->lastConfigIndex;

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

S
Shengliang Guan 已提交
594
  syncNodeRelease(pSyncNode);
595 596 597
  return 0;
}

598
int32_t syncGetSnapshotMetaByIndex(int64_t rid, SyncIndex snapshotIndex, struct SSnapshotMeta* sMeta) {
S
Shengliang Guan 已提交
599
  SSyncNode* pSyncNode = syncNodeAcquire(rid);
600 601 602
  if (pSyncNode == NULL) {
    return -1;
  }
M
Minghao Li 已提交
603
  ASSERT(rid == pSyncNode->rid);
604 605 606 607 608 609 610 611 612 613 614

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

  for (int i = 0; i < pSyncNode->pRaftCfg->configIndexCount; ++i) {
    if ((pSyncNode->pRaftCfg->configIndexArr)[i] > lastIndex &&
        (pSyncNode->pRaftCfg->configIndexArr)[i] <= snapshotIndex) {
      lastIndex = (pSyncNode->pRaftCfg->configIndexArr)[i];
    }
  }
  sMeta->lastConfigIndex = lastIndex;
615
  sTrace("vgId:%d, get snapshot meta by index:%" PRId64 " lcindex:%" PRId64, pSyncNode->vgId, snapshotIndex,
S
Shengliang Guan 已提交
616
         sMeta->lastConfigIndex);
617

S
Shengliang Guan 已提交
618
  syncNodeRelease(pSyncNode);
619 620 621
  return 0;
}

622 623 624 625 626 627 628 629 630 631
SyncIndex syncNodeGetSnapshotConfigIndex(SSyncNode* pSyncNode, SyncIndex snapshotLastApplyIndex) {
  ASSERT(pSyncNode->pRaftCfg->configIndexCount >= 1);
  SyncIndex lastIndex = (pSyncNode->pRaftCfg->configIndexArr)[0];

  for (int i = 0; i < pSyncNode->pRaftCfg->configIndexCount; ++i) {
    if ((pSyncNode->pRaftCfg->configIndexArr)[i] > lastIndex &&
        (pSyncNode->pRaftCfg->configIndexArr)[i] <= snapshotLastApplyIndex) {
      lastIndex = (pSyncNode->pRaftCfg->configIndexArr)[i];
    }
  }
S
Shengliang Guan 已提交
632
  sTrace("vgId:%d, sync get last config index, index:%" PRId64 " lcindex:%" PRId64, pSyncNode->vgId,
S
Shengliang Guan 已提交
633
         snapshotLastApplyIndex, lastIndex);
634 635 636 637

  return lastIndex;
}

M
Minghao Li 已提交
638 639 640 641 642
const char* syncGetMyRoleStr(int64_t rid) {
  const char* s = syncUtilState2String(syncGetMyRole(rid));
  return s;
}

M
Minghao Li 已提交
643
bool syncRestoreFinish(int64_t rid) {
S
Shengliang Guan 已提交
644
  SSyncNode* pSyncNode = syncNodeAcquire(rid);
M
Minghao Li 已提交
645 646 647 648 649 650
  if (pSyncNode == NULL) {
    return false;
  }
  ASSERT(rid == pSyncNode->rid);
  bool restoreFinish = pSyncNode->restoreFinish;

S
Shengliang Guan 已提交
651
  syncNodeRelease(pSyncNode);
M
Minghao Li 已提交
652 653 654
  return restoreFinish;
}

M
Minghao Li 已提交
655
SyncTerm syncGetMyTerm(int64_t rid) {
S
Shengliang Guan 已提交
656
  SSyncNode* pSyncNode = syncNodeAcquire(rid);
M
Minghao Li 已提交
657 658 659
  if (pSyncNode == NULL) {
    return TAOS_SYNC_STATE_ERROR;
  }
M
Minghao Li 已提交
660
  ASSERT(rid == pSyncNode->rid);
M
Minghao Li 已提交
661
  SyncTerm term = pSyncNode->pRaftStore->currentTerm;
M
Minghao Li 已提交
662

S
Shengliang Guan 已提交
663
  syncNodeRelease(pSyncNode);
M
Minghao Li 已提交
664
  return term;
M
Minghao Li 已提交
665 666
}

667
SyncIndex syncGetLastIndex(int64_t rid) {
S
Shengliang Guan 已提交
668
  SSyncNode* pSyncNode = syncNodeAcquire(rid);
669 670 671 672 673 674
  if (pSyncNode == NULL) {
    return SYNC_INDEX_INVALID;
  }
  ASSERT(rid == pSyncNode->rid);
  SyncIndex lastIndex = syncNodeGetLastIndex(pSyncNode);

S
Shengliang Guan 已提交
675
  syncNodeRelease(pSyncNode);
676 677 678 679
  return lastIndex;
}

SyncIndex syncGetCommitIndex(int64_t rid) {
S
Shengliang Guan 已提交
680
  SSyncNode* pSyncNode = syncNodeAcquire(rid);
681 682 683 684 685 686
  if (pSyncNode == NULL) {
    return SYNC_INDEX_INVALID;
  }
  ASSERT(rid == pSyncNode->rid);
  SyncIndex cmtIndex = pSyncNode->commitIndex;

S
Shengliang Guan 已提交
687
  syncNodeRelease(pSyncNode);
688 689 690
  return cmtIndex;
}

M
Minghao Li 已提交
691
SyncGroupId syncGetVgId(int64_t rid) {
S
Shengliang Guan 已提交
692
  SSyncNode* pSyncNode = syncNodeAcquire(rid);
M
Minghao Li 已提交
693
  if (pSyncNode == NULL) {
M
Minghao Li 已提交
694 695
    return TAOS_SYNC_STATE_ERROR;
  }
M
Minghao Li 已提交
696
  ASSERT(rid == pSyncNode->rid);
M
Minghao Li 已提交
697
  SyncGroupId vgId = pSyncNode->vgId;
M
Minghao Li 已提交
698

S
Shengliang Guan 已提交
699
  syncNodeRelease(pSyncNode);
M
Minghao Li 已提交
700
  return vgId;
M
Minghao Li 已提交
701 702
}

M
Minghao Li 已提交
703
void syncGetEpSet(int64_t rid, SEpSet* pEpSet) {
S
Shengliang Guan 已提交
704
  SSyncNode* pSyncNode = syncNodeAcquire(rid);
M
Minghao Li 已提交
705 706 707 708
  if (pSyncNode == NULL) {
    memset(pEpSet, 0, sizeof(*pEpSet));
    return;
  }
M
Minghao Li 已提交
709
  ASSERT(rid == pSyncNode->rid);
M
Minghao Li 已提交
710 711
  pEpSet->numOfEps = 0;
  for (int i = 0; i < pSyncNode->pRaftCfg->cfg.replicaNum; ++i) {
M
Minghao Li 已提交
712 713
    snprintf(pEpSet->eps[i].fqdn, sizeof(pEpSet->eps[i].fqdn), "%s", (pSyncNode->pRaftCfg->cfg.nodeInfo)[i].nodeFqdn);
    pEpSet->eps[i].port = (pSyncNode->pRaftCfg->cfg.nodeInfo)[i].nodePort;
M
Minghao Li 已提交
714
    (pEpSet->numOfEps)++;
S
Shengliang Guan 已提交
715
    sInfo("vgId:%d, sync get epset: index:%d %s:%d", pSyncNode->vgId, i, pEpSet->eps[i].fqdn, pEpSet->eps[i].port);
M
Minghao Li 已提交
716 717
  }
  pEpSet->inUse = pSyncNode->pRaftCfg->cfg.myIndex;
S
Shengliang Guan 已提交
718
  sInfo("vgId:%d, sync get epset in-use:%d", pSyncNode->vgId, pEpSet->inUse);
719

S
Shengliang Guan 已提交
720
  syncNodeRelease(pSyncNode);
721
}
M
Minghao Li 已提交
722

723
void syncGetRetryEpSet(int64_t rid, SEpSet* pEpSet) {
S
Shengliang Guan 已提交
724
  SSyncNode* pSyncNode = syncNodeAcquire(rid);
725 726 727 728 729 730 731 732 733 734
  if (pSyncNode == NULL) {
    memset(pEpSet, 0, sizeof(*pEpSet));
    return;
  }
  ASSERT(rid == pSyncNode->rid);
  pEpSet->numOfEps = 0;
  for (int i = 0; i < pSyncNode->pRaftCfg->cfg.replicaNum; ++i) {
    snprintf(pEpSet->eps[i].fqdn, sizeof(pEpSet->eps[i].fqdn), "%s", (pSyncNode->pRaftCfg->cfg.nodeInfo)[i].nodeFqdn);
    pEpSet->eps[i].port = (pSyncNode->pRaftCfg->cfg.nodeInfo)[i].nodePort;
    (pEpSet->numOfEps)++;
M
Minghao Li 已提交
735 736
    sInfo("vgId:%d, sync get retry epset: index:%d %s:%d", pSyncNode->vgId, i, pEpSet->eps[i].fqdn,
          pEpSet->eps[i].port);
737
  }
M
Minghao Li 已提交
738 739 740
  if (pEpSet->numOfEps > 0) {
    pEpSet->inUse = (pSyncNode->pRaftCfg->cfg.myIndex + 1) % pEpSet->numOfEps;
  }
M
Minghao Li 已提交
741
  sInfo("vgId:%d, sync get retry epset in-use:%d", pSyncNode->vgId, pEpSet->inUse);
M
Minghao Li 已提交
742

S
Shengliang Guan 已提交
743
  syncNodeRelease(pSyncNode);
M
Minghao Li 已提交
744 745
}

S
Shengliang Guan 已提交
746
static void syncGetAndDelRespRpc(SSyncNode* pSyncNode, uint64_t index, SRpcHandleInfo* pInfo) {
M
Minghao Li 已提交
747 748 749
  SRespStub stub;
  int32_t   ret = syncRespMgrGetAndDel(pSyncNode->pSyncRespMgr, index, &stub);
  if (ret == 1) {
S
Shengliang Guan 已提交
750
    *pInfo = stub.rpcMsg.info;
M
Minghao Li 已提交
751 752
  }

S
Shengliang Guan 已提交
753
  sTrace("vgId:%d, get seq:%" PRIu64 " rpc handle:%p", pSyncNode->vgId, index, pInfo->handle);
M
Minghao Li 已提交
754 755 756
}

char* sync2SimpleStr(int64_t rid) {
S
Shengliang Guan 已提交
757
  SSyncNode* pSyncNode = syncNodeAcquire(rid);
M
Minghao Li 已提交
758
  if (pSyncNode == NULL) {
S
Shengliang Guan 已提交
759
    sTrace("syncSetRpc get pSyncNode is NULL, rid:%" PRId64, rid);
M
Minghao Li 已提交
760 761
    return NULL;
  }
M
Minghao Li 已提交
762
  ASSERT(rid == pSyncNode->rid);
M
Minghao Li 已提交
763
  char* s = syncNode2SimpleStr(pSyncNode);
S
Shengliang Guan 已提交
764
  syncNodeRelease(pSyncNode);
M
Minghao Li 已提交
765 766 767 768

  return s;
}

M
Minghao Li 已提交
769
int32_t syncPropose(int64_t rid, SRpcMsg* pMsg, bool isWeak) {
S
Shengliang Guan 已提交
770
  SSyncNode* pSyncNode = syncNodeAcquire(rid);
771
  if (pSyncNode == NULL) {
S
Shengliang Guan 已提交
772
    syncNodeRelease(pSyncNode);
M
Minghao Li 已提交
773 774
    terrno = TSDB_CODE_SYN_INTERNAL_ERROR;
    return -1;
775
  }
M
Minghao Li 已提交
776
  ASSERT(rid == pSyncNode->rid);
M
Minghao Li 已提交
777

778
  int32_t ret = syncNodePropose(pSyncNode, pMsg, isWeak);
S
Shengliang Guan 已提交
779
  syncNodeRelease(pSyncNode);
M
Minghao Li 已提交
780 781 782
  return ret;
}

783
static bool syncNodeBatchOK(SRpcMsg** pMsgPArr, int32_t arrSize) {
M
Minghao Li 已提交
784
  for (int32_t i = 0; i < arrSize; ++i) {
785
    if (pMsgPArr[i]->msgType == TDMT_SYNC_CONFIG_CHANGE) {
M
Minghao Li 已提交
786 787 788
      return false;
    }

789
    if (pMsgPArr[i]->msgType == TDMT_SYNC_CONFIG_CHANGE_FINISH) {
M
Minghao Li 已提交
790 791 792 793 794 795 796
      return false;
    }
  }

  return true;
}

797
int32_t syncNodePropose(SSyncNode* pSyncNode, SRpcMsg* pMsg, bool isWeak) {
M
Minghao Li 已提交
798
  int32_t ret = 0;
M
Minghao Li 已提交
799

M
Minghao Li 已提交
800 801
  do {
    char eventLog[128];
S
Shengliang Guan 已提交
802
    snprintf(eventLog, sizeof(eventLog), "propose message, type:%s", TMSG_INFO(pMsg->msgType));
M
Minghao Li 已提交
803 804
    syncNodeEventLog(pSyncNode, eventLog);
  } while (0);
M
Minghao Li 已提交
805

M
Minghao Li 已提交
806
  if (pSyncNode->state == TAOS_SYNC_STATE_LEADER) {
M
Minghao Li 已提交
807 808 809
    if (pSyncNode->changing && pMsg->msgType != TDMT_SYNC_CONFIG_CHANGE_FINISH) {
      ret = -1;
      terrno = TSDB_CODE_SYN_PROPOSE_NOT_READY;
S
Shengliang Guan 已提交
810
      sError("vgId:%d, failed to sync propose since not ready, type:%s", pSyncNode->vgId, TMSG_INFO(pMsg->msgType));
M
Minghao Li 已提交
811 812 813 814 815 816 817 818
      goto _END;
    }

    // config change
    if (pMsg->msgType == TDMT_SYNC_CONFIG_CHANGE) {
      if (!syncNodeCanChange(pSyncNode)) {
        ret = -1;
        terrno = TSDB_CODE_SYN_RECONFIG_NOT_READY;
S
Shengliang Guan 已提交
819
        sError("vgId:%d, failed to sync reconfig since not ready, type:%s", pSyncNode->vgId, TMSG_INFO(pMsg->msgType));
M
Minghao Li 已提交
820 821 822 823 824 825 826
        goto _END;
      }

      ASSERT(!pSyncNode->changing);
      pSyncNode->changing = true;
    }

827 828
    // not restored, vnode enable
    if (!pSyncNode->restoreFinish && pSyncNode->vgId != 1) {
829 830
      ret = -1;
      terrno = TSDB_CODE_SYN_PROPOSE_NOT_READY;
S
Shengliang Guan 已提交
831 832
      sError("vgId:%d, failed to sync propose since not ready, type:%s, last:%" PRId64 ", cmt:%" PRId64,
             pSyncNode->vgId, TMSG_INFO(pMsg->msgType), syncNodeGetLastIndex(pSyncNode), pSyncNode->commitIndex);
833 834 835
      goto _END;
    }

M
Minghao Li 已提交
836 837 838 839 840 841
    SRespStub stub;
    stub.createTime = taosGetTimestampMs();
    stub.rpcMsg = *pMsg;
    uint64_t seqNum = syncRespMgrAdd(pSyncNode->pSyncRespMgr, &stub);

    SyncClientRequest* pSyncMsg = syncClientRequestBuild2(pMsg, seqNum, isWeak, pSyncNode->vgId);
M
Minghao Li 已提交
842 843
    SRpcMsg            rpcMsg;
    syncClientRequest2RpcMsg(pSyncMsg, &rpcMsg);
844

845 846 847
    // optimized one replica
    if (syncNodeIsOptimizedOneReplica(pSyncNode, pMsg)) {
      SyncIndex retIndex;
M
Minghao Li 已提交
848
      int32_t   code = syncNodeOnClientRequest(pSyncNode, pSyncMsg, &retIndex);
849 850
      if (code == 0) {
        pMsg->info.conn.applyIndex = retIndex;
M
Minghao Li 已提交
851
        pMsg->info.conn.applyTerm = pSyncNode->pRaftStore->currentTerm;
852 853 854
        rpcFreeCont(rpcMsg.pCont);
        syncRespMgrDel(pSyncNode->pSyncRespMgr, seqNum);
        ret = 1;
855 856
        sDebug("vgId:%d, sync optimize index:%" PRId64 ", type:%s", pSyncNode->vgId, retIndex,
               TMSG_INFO(pMsg->msgType));
857 858 859
      } else {
        ret = -1;
        terrno = TSDB_CODE_SYN_INTERNAL_ERROR;
S
Shengliang Guan 已提交
860 861
        sError("vgId:%d, failed to sync optimize index:%" PRId64 ", type:%s", pSyncNode->vgId, retIndex,
               TMSG_INFO(pMsg->msgType));
862 863
      }

M
Minghao Li 已提交
864
    } else {
S
Shengliang Guan 已提交
865
      if (pSyncNode->syncEqMsg != NULL && (*pSyncNode->syncEqMsg)(pSyncNode->msgcb, &rpcMsg) == 0) {
866 867 868 869
        ret = 0;
      } else {
        ret = -1;
        terrno = TSDB_CODE_SYN_INTERNAL_ERROR;
S
Shengliang Guan 已提交
870
        sError("vgId:%d, failed to enqueue msg since its null", pSyncNode->vgId);
871
      }
M
Minghao Li 已提交
872
    }
873

M
Minghao Li 已提交
874
    syncClientRequestDestroy(pSyncMsg);
M
Minghao Li 已提交
875 876
    goto _END;

M
Minghao Li 已提交
877
  } else {
M
Minghao Li 已提交
878 879
    ret = -1;
    terrno = TSDB_CODE_SYN_NOT_LEADER;
S
Shengliang Guan 已提交
880 881
    sError("vgId:%d, sync propose not leader, %s, type:%s", pSyncNode->vgId, syncUtilState2String(pSyncNode->state),
           TMSG_INFO(pMsg->msgType));
M
Minghao Li 已提交
882
    goto _END;
M
Minghao Li 已提交
883
  }
M
Minghao Li 已提交
884

M
Minghao Li 已提交
885
_END:
M
Minghao Li 已提交
886 887 888
  return ret;
}

889 890 891 892 893 894 895 896 897 898 899 900
int32_t syncHbTimerInit(SSyncNode* pSyncNode, SSyncTimer* pSyncTimer, SRaftId destId) {
  pSyncTimer->pTimer = NULL;
  pSyncTimer->counter = 0;
  pSyncTimer->timerMS = pSyncNode->hbBaseLine;
  pSyncTimer->timerCb = syncNodeEqPeerHeartbeatTimer;
  pSyncTimer->destId = destId;
  atomic_store_64(&pSyncTimer->logicClock, 0);
  return 0;
}

int32_t syncHbTimerStart(SSyncNode* pSyncNode, SSyncTimer* pSyncTimer) {
  int32_t ret = 0;
S
Shengliang Guan 已提交
901
  if (syncIsInit()) {
M
Minghao Li 已提交
902
    SSyncHbTimerData* pData = taosMemoryMalloc(sizeof(SSyncHbTimerData));
903 904 905 906
    pData->pSyncNode = pSyncNode;
    pData->pTimer = pSyncTimer;
    pData->destId = pSyncTimer->destId;
    pData->logicClock = pSyncTimer->logicClock;
M
Minghao Li 已提交
907

908
    pSyncTimer->pData = pData;
S
Shengliang Guan 已提交
909
    taosTmrReset(pSyncTimer->timerCb, pSyncTimer->timerMS, pData, syncEnv()->pTimerManager, &pSyncTimer->pTimer);
910 911 912 913 914 915 916 917 918 919 920
  } else {
    sError("vgId:%d, start ctrl hb timer error, sync env is stop", pSyncNode->vgId);
  }
  return ret;
}

int32_t syncHbTimerStop(SSyncNode* pSyncNode, SSyncTimer* pSyncTimer) {
  int32_t ret = 0;
  atomic_add_fetch_64(&pSyncTimer->logicClock, 1);
  taosTmrStop(pSyncTimer->pTimer);
  pSyncTimer->pTimer = NULL;
M
Minghao Li 已提交
921
  // taosMemoryFree(pSyncTimer->pData);
922 923 924
  return ret;
}

S
Shengliang Guan 已提交
925 926
SSyncNode* syncNodeOpen(SSyncInfo* pSyncInfo) {
  SSyncNode* pSyncNode = taosMemoryCalloc(1, sizeof(SSyncNode));
927 928 929 930
  if (pSyncNode == NULL) {
    terrno = TSDB_CODE_OUT_OF_MEMORY;
    goto _error;
  }
M
Minghao Li 已提交
931

M
Minghao Li 已提交
932 933 934 935
  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());
936
      goto _error;
M
Minghao Li 已提交
937
    }
938
  }
M
Minghao Li 已提交
939

S
Shengliang Guan 已提交
940
  snprintf(pSyncNode->configPath, sizeof(pSyncNode->configPath), "%s%sraft_config.json", pSyncInfo->path, TD_DIRSEP);
941
  if (!taosCheckExistFile(pSyncNode->configPath)) {
M
Minghao Li 已提交
942
    // create a new raft config file
S
Shengliang Guan 已提交
943
    SRaftCfgMeta meta = {0};
M
Minghao Li 已提交
944
    meta.isStandBy = pSyncInfo->isStandBy;
M
Minghao Li 已提交
945
    meta.snapshotStrategy = pSyncInfo->snapshotStrategy;
946
    meta.lastConfigIndex = SYNC_INDEX_INVALID;
M
Minghao Li 已提交
947
    meta.batchSize = pSyncInfo->batchSize;
S
Shengliang Guan 已提交
948 949
    if (raftCfgCreateFile(&pSyncInfo->syncCfg, meta, pSyncNode->configPath) != 0) {
      sError("vgId:%d, failed to create raft cfg file at %s", pSyncNode->vgId, pSyncNode->configPath);
H
Hongze Cheng 已提交
950
      goto _error;
951
    }
952
    if (pSyncInfo->syncCfg.replicaNum == 0) {
S
Shengliang Guan 已提交
953
      sInfo("vgId:%d, sync config not input", pSyncNode->vgId);
954 955
      pSyncInfo->syncCfg = pSyncNode->pRaftCfg->cfg;
    }
956 957 958
  } else {
    // update syncCfg by raft_config.json
    pSyncNode->pRaftCfg = raftCfgOpen(pSyncNode->configPath);
959
    if (pSyncNode->pRaftCfg == NULL) {
S
Shengliang Guan 已提交
960
      sError("vgId:%d, failed to open raft cfg file at %s", pSyncNode->vgId, pSyncNode->configPath);
H
Hongze Cheng 已提交
961
      goto _error;
962
    }
S
Shengliang Guan 已提交
963 964

    if (pSyncInfo->syncCfg.replicaNum > 0 && syncIsConfigChanged(&pSyncNode->pRaftCfg->cfg, &pSyncInfo->syncCfg)) {
S
Shengliang Guan 已提交
965 966 967 968 969 970
      sInfo("vgId:%d, use sync config from input options and write to cfg file", pSyncNode->vgId);
      pSyncNode->pRaftCfg->cfg = pSyncInfo->syncCfg;
      if (raftCfgPersist(pSyncNode->pRaftCfg) != 0) {
        sError("vgId:%d, failed to persist raft cfg file at %s", pSyncNode->vgId, pSyncNode->configPath);
        goto _error;
      }
S
Shengliang Guan 已提交
971 972 973 974
    } else {
      sInfo("vgId:%d, use sync config from raft cfg file", pSyncNode->vgId);
      pSyncInfo->syncCfg = pSyncNode->pRaftCfg->cfg;
    }
975 976

    raftCfgClose(pSyncNode->pRaftCfg);
977
    pSyncNode->pRaftCfg = NULL;
M
Minghao Li 已提交
978 979
  }

S
Shengliang Guan 已提交
980 981
  // init by SSyncInfo
  pSyncNode->vgId = pSyncInfo->vgId;
S
Shengliang Guan 已提交
982 983 984 985 986 987 988
  SSyncCfg* pCfg = &pSyncInfo->syncCfg;
  sDebug("vgId:%d, replica:%d selfIndex:%d", pSyncNode->vgId, pCfg->replicaNum, pCfg->myIndex);
  for (int32_t i = 0; i < pCfg->replicaNum; ++i) {
    SNodeInfo* pNode = &pCfg->nodeInfo[i];
    sDebug("vgId:%d, index:%d ep:%s:%u", pSyncNode->vgId, i, pNode->nodeFqdn, pNode->nodePort);
  }

M
Minghao Li 已提交
989
  memcpy(pSyncNode->path, pSyncInfo->path, sizeof(pSyncNode->path));
S
Shengliang Guan 已提交
990 991 992
  snprintf(pSyncNode->raftStorePath, sizeof(pSyncNode->raftStorePath), "%s%sraft_store.json", pSyncInfo->path,
           TD_DIRSEP);
  snprintf(pSyncNode->configPath, sizeof(pSyncNode->configPath), "%s%sraft_config.json", pSyncInfo->path, TD_DIRSEP);
M
Minghao Li 已提交
993

M
Minghao Li 已提交
994
  pSyncNode->pWal = pSyncInfo->pWal;
S
Shengliang Guan 已提交
995
  pSyncNode->msgcb = pSyncInfo->msgcb;
S
Shengliang Guan 已提交
996 997 998
  pSyncNode->syncSendMSg = pSyncInfo->syncSendMSg;
  pSyncNode->syncEqMsg = pSyncInfo->syncEqMsg;
  pSyncNode->syncEqCtrlMsg = pSyncInfo->syncEqCtrlMsg;
M
Minghao Li 已提交
999

M
Minghao Li 已提交
1000 1001
  // init raft config
  pSyncNode->pRaftCfg = raftCfgOpen(pSyncNode->configPath);
1002
  if (pSyncNode->pRaftCfg == NULL) {
S
Shengliang Guan 已提交
1003
    sError("vgId:%d, failed to open raft cfg file at %s", pSyncNode->vgId, pSyncNode->configPath);
1004 1005
    goto _error;
  }
M
Minghao Li 已提交
1006

M
Minghao Li 已提交
1007
  // init internal
M
Minghao Li 已提交
1008
  pSyncNode->myNodeInfo = pSyncNode->pRaftCfg->cfg.nodeInfo[pSyncNode->pRaftCfg->cfg.myIndex];
1009
  if (!syncUtilnodeInfo2raftId(&pSyncNode->myNodeInfo, pSyncNode->vgId, &pSyncNode->myRaftId)) {
S
Shengliang Guan 已提交
1010
    sError("vgId:%d, failed to determine my raft member id", pSyncNode->vgId);
H
Hongze Cheng 已提交
1011
    goto _error;
1012
  }
M
Minghao Li 已提交
1013

M
Minghao Li 已提交
1014
  // init peersNum, peers, peersId
M
Minghao Li 已提交
1015
  pSyncNode->peersNum = pSyncNode->pRaftCfg->cfg.replicaNum - 1;
M
Minghao Li 已提交
1016
  int j = 0;
M
Minghao Li 已提交
1017 1018 1019
  for (int i = 0; i < pSyncNode->pRaftCfg->cfg.replicaNum; ++i) {
    if (i != pSyncNode->pRaftCfg->cfg.myIndex) {
      pSyncNode->peersNodeInfo[j] = pSyncNode->pRaftCfg->cfg.nodeInfo[i];
M
Minghao Li 已提交
1020 1021 1022
      j++;
    }
  }
M
Minghao Li 已提交
1023
  for (int i = 0; i < pSyncNode->peersNum; ++i) {
1024
    if (!syncUtilnodeInfo2raftId(&pSyncNode->peersNodeInfo[i], pSyncNode->vgId, &pSyncNode->peersId[i])) {
S
Shengliang Guan 已提交
1025
      sError("vgId:%d, failed to determine raft member id, peer:%d", pSyncNode->vgId, i);
H
Hongze Cheng 已提交
1026
      goto _error;
1027
    }
M
Minghao Li 已提交
1028
  }
M
Minghao Li 已提交
1029

M
Minghao Li 已提交
1030
  // init replicaNum, replicasId
M
Minghao Li 已提交
1031 1032
  pSyncNode->replicaNum = pSyncNode->pRaftCfg->cfg.replicaNum;
  for (int i = 0; i < pSyncNode->pRaftCfg->cfg.replicaNum; ++i) {
H
Hongze Cheng 已提交
1033
    if (!syncUtilnodeInfo2raftId(&pSyncNode->pRaftCfg->cfg.nodeInfo[i], pSyncNode->vgId, &pSyncNode->replicasId[i])) {
S
Shengliang Guan 已提交
1034
      sError("vgId:%d, failed to determine raft member id, replica:%d", pSyncNode->vgId, i);
H
Hongze Cheng 已提交
1035
      goto _error;
1036
    }
M
Minghao Li 已提交
1037 1038
  }

M
Minghao Li 已提交
1039
  // init raft algorithm
M
Minghao Li 已提交
1040
  pSyncNode->pFsm = pSyncInfo->pFsm;
1041
  pSyncInfo->pFsm = NULL;
M
Minghao Li 已提交
1042
  pSyncNode->quorum = syncUtilQuorum(pSyncNode->pRaftCfg->cfg.replicaNum);
M
Minghao Li 已提交
1043 1044
  pSyncNode->leaderCache = EMPTY_RAFT_ID;

M
Minghao Li 已提交
1045
  // init life cycle outside
M
Minghao Li 已提交
1046

M
Minghao Li 已提交
1047 1048 1049 1050 1051 1052 1053 1054 1055 1056 1057 1058 1059 1060 1061 1062 1063 1064 1065 1066 1067 1068 1069 1070
  // 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 已提交
1071
  // init TLA+ server vars
M
syncInt  
Minghao Li 已提交
1072
  pSyncNode->state = TAOS_SYNC_STATE_FOLLOWER;
M
Minghao Li 已提交
1073
  pSyncNode->pRaftStore = raftStoreOpen(pSyncNode->raftStorePath);
1074
  if (pSyncNode->pRaftStore == NULL) {
S
Shengliang Guan 已提交
1075
    sError("vgId:%d, failed to open raft store at path %s", pSyncNode->vgId, pSyncNode->raftStorePath);
1076 1077
    goto _error;
  }
M
Minghao Li 已提交
1078

M
Minghao Li 已提交
1079
  // init TLA+ candidate vars
M
Minghao Li 已提交
1080
  pSyncNode->pVotesGranted = voteGrantedCreate(pSyncNode);
1081
  if (pSyncNode->pVotesGranted == NULL) {
S
Shengliang Guan 已提交
1082
    sError("vgId:%d, failed to create VotesGranted", pSyncNode->vgId);
1083 1084
    goto _error;
  }
M
Minghao Li 已提交
1085
  pSyncNode->pVotesRespond = votesRespondCreate(pSyncNode);
1086
  if (pSyncNode->pVotesRespond == NULL) {
S
Shengliang Guan 已提交
1087
    sError("vgId:%d, failed to create VotesRespond", pSyncNode->vgId);
1088 1089
    goto _error;
  }
M
Minghao Li 已提交
1090

M
Minghao Li 已提交
1091 1092
  // init TLA+ leader vars
  pSyncNode->pNextIndex = syncIndexMgrCreate(pSyncNode);
1093
  if (pSyncNode->pNextIndex == NULL) {
S
Shengliang Guan 已提交
1094
    sError("vgId:%d, failed to create SyncIndexMgr", pSyncNode->vgId);
1095 1096
    goto _error;
  }
M
Minghao Li 已提交
1097
  pSyncNode->pMatchIndex = syncIndexMgrCreate(pSyncNode);
1098
  if (pSyncNode->pMatchIndex == NULL) {
S
Shengliang Guan 已提交
1099
    sError("vgId:%d, failed to create SyncIndexMgr", pSyncNode->vgId);
1100 1101
    goto _error;
  }
M
Minghao Li 已提交
1102 1103 1104

  // init TLA+ log vars
  pSyncNode->pLogStore = logStoreCreate(pSyncNode);
1105
  if (pSyncNode->pLogStore == NULL) {
S
Shengliang Guan 已提交
1106
    sError("vgId:%d, failed to create SyncLogStore", pSyncNode->vgId);
1107 1108
    goto _error;
  }
1109 1110 1111 1112 1113

  SyncIndex commitIndex = SYNC_INDEX_INVALID;
  if (pSyncNode->pFsm != NULL && pSyncNode->pFsm->FpGetSnapshotInfo != NULL) {
    SSnapshot snapshot = {0};
    int32_t   code = pSyncNode->pFsm->FpGetSnapshotInfo(pSyncNode->pFsm, &snapshot);
1114
    if (code != 0) {
S
Shengliang Guan 已提交
1115
      sError("vgId:%d, failed to get snapshot info, code:%d", pSyncNode->vgId, code);
H
Hongze Cheng 已提交
1116
      goto _error;
1117
    }
1118 1119 1120 1121 1122 1123
    if (snapshot.lastApplyIndex > commitIndex) {
      commitIndex = snapshot.lastApplyIndex;
      syncNodeEventLog(pSyncNode, "reset commit index by snapshot");
    }
  }
  pSyncNode->commitIndex = commitIndex;
M
Minghao Li 已提交
1124

M
Minghao Li 已提交
1125 1126 1127 1128 1129
  // timer ms init
  pSyncNode->pingBaseLine = PING_TIMER_MS;
  pSyncNode->electBaseLine = ELECT_TIMER_MS_MIN;
  pSyncNode->hbBaseLine = HEARTBEAT_TIMER_MS;

M
Minghao Li 已提交
1130
  // init ping timer
M
Minghao Li 已提交
1131
  pSyncNode->pPingTimer = NULL;
M
Minghao Li 已提交
1132
  pSyncNode->pingTimerMS = pSyncNode->pingBaseLine;
M
Minghao Li 已提交
1133 1134
  atomic_store_64(&pSyncNode->pingTimerLogicClock, 0);
  atomic_store_64(&pSyncNode->pingTimerLogicClockUser, 0);
M
Minghao Li 已提交
1135
  pSyncNode->FpPingTimerCB = syncNodeEqPingTimer;
M
Minghao Li 已提交
1136
  pSyncNode->pingTimerCounter = 0;
M
Minghao Li 已提交
1137

M
Minghao Li 已提交
1138 1139
  // init elect timer
  pSyncNode->pElectTimer = NULL;
M
Minghao Li 已提交
1140
  pSyncNode->electTimerMS = syncUtilElectRandomMS(pSyncNode->electBaseLine, 2 * pSyncNode->electBaseLine);
M
Minghao Li 已提交
1141
  atomic_store_64(&pSyncNode->electTimerLogicClock, 0);
M
Minghao Li 已提交
1142
  pSyncNode->FpElectTimerCB = syncNodeEqElectTimer;
M
Minghao Li 已提交
1143 1144 1145 1146
  pSyncNode->electTimerCounter = 0;

  // init heartbeat timer
  pSyncNode->pHeartbeatTimer = NULL;
M
Minghao Li 已提交
1147
  pSyncNode->heartbeatTimerMS = pSyncNode->hbBaseLine;
M
Minghao Li 已提交
1148 1149
  atomic_store_64(&pSyncNode->heartbeatTimerLogicClock, 0);
  atomic_store_64(&pSyncNode->heartbeatTimerLogicClockUser, 0);
M
Minghao Li 已提交
1150
  pSyncNode->FpHeartbeatTimerCB = syncNodeEqHeartbeatTimer;
M
Minghao Li 已提交
1151 1152
  pSyncNode->heartbeatTimerCounter = 0;

1153 1154 1155 1156 1157
  // 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 已提交
1158
  // init callback
M
Minghao Li 已提交
1159 1160
  pSyncNode->FpOnPing = syncNodeOnPing;
  pSyncNode->FpOnPingReply = syncNodeOnPingReply;
M
Minghao Li 已提交
1161
  pSyncNode->FpOnClientRequest = syncNodeOnClientRequest;
1162
  pSyncNode->FpOnTimeout = syncNodeOnTimer;
M
Minghao Li 已提交
1163 1164 1165 1166 1167 1168
  pSyncNode->FpOnSnapshot = syncNodeOnSnapshot;
  pSyncNode->FpOnSnapshotReply = syncNodeOnSnapshotReply;
  pSyncNode->FpOnRequestVote = syncNodeOnRequestVote;
  pSyncNode->FpOnRequestVoteReply = syncNodeOnRequestVoteReply;
  pSyncNode->FpOnAppendEntries = syncNodeOnAppendEntries;
  pSyncNode->FpOnAppendEntriesReply = syncNodeOnAppendEntriesReply;
M
Minghao Li 已提交
1169

M
Minghao Li 已提交
1170
  // tools
M
Minghao Li 已提交
1171
  pSyncNode->pSyncRespMgr = syncRespMgrCreate(pSyncNode, SYNC_RESP_TTL_MS);
1172
  if (pSyncNode->pSyncRespMgr == NULL) {
S
Shengliang Guan 已提交
1173
    sError("vgId:%d, failed to create SyncRespMgr", pSyncNode->vgId);
1174 1175
    goto _error;
  }
M
Minghao Li 已提交
1176

1177 1178
  // restore state
  pSyncNode->restoreFinish = false;
1179

M
Minghao Li 已提交
1180 1181 1182 1183 1184 1185 1186 1187
  // snapshot senders
  for (int i = 0; i < TSDB_MAX_REPLICA; ++i) {
    SSyncSnapshotSender* pSender = snapshotSenderCreate(pSyncNode, i);
    // ASSERT(pSender != NULL);
    (pSyncNode->senders)[i] = pSender;
  }

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

M
Minghao Li 已提交
1190 1191 1192
  // is config changing
  pSyncNode->changing = false;

M
Minghao Li 已提交
1193 1194 1195
  // peer state
  syncNodePeerStateInit(pSyncNode);

M
Minghao Li 已提交
1196 1197 1198
  // min match index
  pSyncNode->minMatchIndex = SYNC_INDEX_INVALID;

M
Minghao Li 已提交
1199
  // start in syncNodeStart
M
Minghao Li 已提交
1200
  // start raft
M
Minghao Li 已提交
1201
  // syncNodeBecomeFollower(pSyncNode);
M
Minghao Li 已提交
1202

M
Minghao Li 已提交
1203 1204
  int64_t timeNow = taosGetTimestampMs();
  pSyncNode->startTime = timeNow;
1205
  pSyncNode->leaderTime = timeNow;
M
Minghao Li 已提交
1206 1207
  pSyncNode->lastReplicateTime = timeNow;

1208 1209 1210
  // snapshotting
  atomic_store_64(&pSyncNode->snapshottingIndex, SYNC_INDEX_INVALID);

M
Minghao Li 已提交
1211
  syncNodeEventLog(pSyncNode, "sync open");
1212

M
Minghao Li 已提交
1213
  return pSyncNode;
1214 1215 1216

_error:
  if (pSyncInfo->pFsm) {
H
Hongze Cheng 已提交
1217 1218
    taosMemoryFree(pSyncInfo->pFsm);
    pSyncInfo->pFsm = NULL;
1219 1220 1221 1222
  }
  syncNodeClose(pSyncNode);
  pSyncNode = NULL;
  return NULL;
M
Minghao Li 已提交
1223 1224
}

M
Minghao Li 已提交
1225 1226 1227 1228 1229 1230 1231 1232 1233 1234 1235
void syncNodeMaybeUpdateCommitBySnapshot(SSyncNode* pSyncNode) {
  if (pSyncNode->pFsm != NULL && pSyncNode->pFsm->FpGetSnapshotInfo != NULL) {
    SSnapshot snapshot;
    int32_t   code = pSyncNode->pFsm->FpGetSnapshotInfo(pSyncNode->pFsm, &snapshot);
    ASSERT(code == 0);
    if (snapshot.lastApplyIndex > pSyncNode->commitIndex) {
      pSyncNode->commitIndex = snapshot.lastApplyIndex;
    }
  }
}

M
Minghao Li 已提交
1236 1237
void syncNodeStart(SSyncNode* pSyncNode) {
  // start raft
1238
  if (pSyncNode->replicaNum == 1) {
M
Minghao Li 已提交
1239
    raftStoreNextTerm(pSyncNode->pRaftStore);
1240
    syncNodeBecomeLeader(pSyncNode, "one replica start");
M
format  
Minghao Li 已提交
1241

1242
    // Raft 3.6.2 Committing entries from previous terms
1243 1244
    syncNodeAppendNoop(pSyncNode);
    syncMaybeAdvanceCommitIndex(pSyncNode);
M
Minghao Li 已提交
1245

M
Minghao Li 已提交
1246 1247
  } else {
    syncNodeBecomeFollower(pSyncNode, "first start");
1248 1249
  }

1250 1251 1252
  int32_t ret = 0;
  ret = syncNodeStartPingTimer(pSyncNode);
  ASSERT(ret == 0);
M
Minghao Li 已提交
1253 1254
}

M
Minghao Li 已提交
1255 1256 1257 1258 1259 1260 1261 1262 1263
void syncNodeStartStandBy(SSyncNode* pSyncNode) {
  // state change
  pSyncNode->state = TAOS_SYNC_STATE_FOLLOWER;
  syncNodeStopHeartbeatTimer(pSyncNode);

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

1265 1266 1267
  ret = 0;
  ret = syncNodeStartPingTimer(pSyncNode);
  ASSERT(ret == 0);
M
Minghao Li 已提交
1268 1269
}

M
Minghao Li 已提交
1270
void syncNodeClose(SSyncNode* pSyncNode) {
1271 1272 1273
  if (pSyncNode == NULL) {
    return;
  }
M
Minghao Li 已提交
1274 1275
  int32_t ret;

M
Minghao Li 已提交
1276 1277
  syncNodeEventLog(pSyncNode, "sync close");

M
Minghao Li 已提交
1278
  ret = raftStoreClose(pSyncNode->pRaftStore);
M
Minghao Li 已提交
1279
  ASSERT(ret == 0);
M
Minghao Li 已提交
1280

M
Minghao Li 已提交
1281
  syncRespMgrDestroy(pSyncNode->pSyncRespMgr);
1282
  pSyncNode->pSyncRespMgr = NULL;
M
Minghao Li 已提交
1283
  voteGrantedDestroy(pSyncNode->pVotesGranted);
1284
  pSyncNode->pVotesGranted = NULL;
M
Minghao Li 已提交
1285
  votesRespondDestory(pSyncNode->pVotesRespond);
1286
  pSyncNode->pVotesRespond = NULL;
M
Minghao Li 已提交
1287
  syncIndexMgrDestroy(pSyncNode->pNextIndex);
1288
  pSyncNode->pNextIndex = NULL;
M
Minghao Li 已提交
1289
  syncIndexMgrDestroy(pSyncNode->pMatchIndex);
1290
  pSyncNode->pMatchIndex = NULL;
M
Minghao Li 已提交
1291
  logStoreDestory(pSyncNode->pLogStore);
1292
  pSyncNode->pLogStore = NULL;
M
Minghao Li 已提交
1293
  raftCfgClose(pSyncNode->pRaftCfg);
1294
  pSyncNode->pRaftCfg = NULL;
M
Minghao Li 已提交
1295 1296 1297 1298 1299

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

M
Minghao Li 已提交
1300 1301 1302 1303
  if (pSyncNode->pFsm != NULL) {
    taosMemoryFree(pSyncNode->pFsm);
  }

M
Minghao Li 已提交
1304 1305 1306 1307 1308 1309 1310
  for (int i = 0; i < TSDB_MAX_REPLICA; ++i) {
    if ((pSyncNode->senders)[i] != NULL) {
      snapshotSenderDestroy((pSyncNode->senders)[i]);
      (pSyncNode->senders)[i] = NULL;
    }
  }

M
Minghao Li 已提交
1311 1312 1313 1314 1315
  if (pSyncNode->pNewNodeReceiver != NULL) {
    snapshotReceiverDestroy(pSyncNode->pNewNodeReceiver);
    pSyncNode->pNewNodeReceiver = NULL;
  }

1316
  taosMemoryFree(pSyncNode);
M
Minghao Li 已提交
1317 1318
}

M
Minghao Li 已提交
1319
// option
M
Minghao Li 已提交
1320 1321
// bool syncNodeSnapshotEnable(SSyncNode* pSyncNode) { return pSyncNode->pRaftCfg->snapshotEnable; }

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

M
Minghao Li 已提交
1324 1325 1326 1327 1328 1329 1330 1331 1332 1333 1334 1335 1336 1337 1338
// ping --------------
int32_t syncNodePing(SSyncNode* pSyncNode, const SRaftId* destRaftId, SyncPing* pMsg) {
  syncPingLog2((char*)"==syncNodePing==", pMsg);
  int32_t ret = 0;

  SRpcMsg rpcMsg;
  syncPing2RpcMsg(pMsg, &rpcMsg);
  syncRpcMsgLog2((char*)"==syncNodePing==", &rpcMsg);

  ret = syncNodeSendMsgById(destRaftId, pSyncNode, &rpcMsg);
  return ret;
}

int32_t syncNodePingSelf(SSyncNode* pSyncNode) {
  int32_t   ret = 0;
M
Minghao Li 已提交
1339
  SyncPing* pMsg = syncPingBuild3(&pSyncNode->myRaftId, &pSyncNode->myRaftId, pSyncNode->vgId);
M
Minghao Li 已提交
1340
  ret = syncNodePing(pSyncNode, &pMsg->destId, pMsg);
M
Minghao Li 已提交
1341
  ASSERT(ret == 0);
M
Minghao Li 已提交
1342 1343 1344 1345 1346 1347 1348 1349

  syncPingDestroy(pMsg);
  return ret;
}

int32_t syncNodePingPeers(SSyncNode* pSyncNode) {
  int32_t ret = 0;
  for (int i = 0; i < pSyncNode->peersNum; ++i) {
M
Minghao Li 已提交
1350 1351 1352
    SRaftId*  destId = &(pSyncNode->peersId[i]);
    SyncPing* pMsg = syncPingBuild3(&pSyncNode->myRaftId, destId, pSyncNode->vgId);
    ret = syncNodePing(pSyncNode, destId, pMsg);
M
Minghao Li 已提交
1353
    ASSERT(ret == 0);
M
Minghao Li 已提交
1354 1355 1356 1357 1358 1359 1360
    syncPingDestroy(pMsg);
  }
  return ret;
}

int32_t syncNodePingAll(SSyncNode* pSyncNode) {
  int32_t ret = 0;
M
Minghao Li 已提交
1361 1362 1363 1364
  for (int i = 0; i < pSyncNode->pRaftCfg->cfg.replicaNum; ++i) {
    SRaftId*  destId = &(pSyncNode->replicasId[i]);
    SyncPing* pMsg = syncPingBuild3(&pSyncNode->myRaftId, destId, pSyncNode->vgId);
    ret = syncNodePing(pSyncNode, destId, pMsg);
M
Minghao Li 已提交
1365
    ASSERT(ret == 0);
M
Minghao Li 已提交
1366 1367 1368 1369 1370 1371 1372 1373
    syncPingDestroy(pMsg);
  }
  return ret;
}

// timer control --------------
int32_t syncNodeStartPingTimer(SSyncNode* pSyncNode) {
  int32_t ret = 0;
S
Shengliang Guan 已提交
1374 1375
  if (syncIsInit()) {
    taosTmrReset(pSyncNode->FpPingTimerCB, pSyncNode->pingTimerMS, pSyncNode, syncEnv()->pTimerManager,
1376 1377 1378
                 &pSyncNode->pPingTimer);
    atomic_store_64(&pSyncNode->pingTimerLogicClock, pSyncNode->pingTimerLogicClockUser);
  } else {
M
Minghao Li 已提交
1379
    sError("vgId:%d, start ping timer error, sync env is stop", pSyncNode->vgId);
1380
  }
M
Minghao Li 已提交
1381 1382 1383 1384 1385 1386 1387 1388 1389 1390 1391 1392 1393
  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 已提交
1394
  if (syncIsInit()) {
1395
    pSyncNode->electTimerMS = ms;
M
Minghao Li 已提交
1396 1397 1398 1399 1400 1401

    SElectTimer* pElectTimer = taosMemoryMalloc(sizeof(SElectTimer));
    pElectTimer->logicClock = pSyncNode->electTimerLogicClock;
    pElectTimer->pSyncNode = pSyncNode;
    pElectTimer->pData = NULL;

S
Shengliang Guan 已提交
1402
    taosTmrReset(pSyncNode->FpElectTimerCB, pSyncNode->electTimerMS, pElectTimer, syncEnv()->pTimerManager,
1403
                 &pSyncNode->pElectTimer);
1404

1405
  } else {
M
Minghao Li 已提交
1406
    sError("vgId:%d, start elect timer error, sync env is stop", pSyncNode->vgId);
1407
  }
M
Minghao Li 已提交
1408 1409 1410 1411 1412
  return ret;
}

int32_t syncNodeStopElectTimer(SSyncNode* pSyncNode) {
  int32_t ret = 0;
M
Minghao Li 已提交
1413
  atomic_add_fetch_64(&pSyncNode->electTimerLogicClock, 1);
M
Minghao Li 已提交
1414 1415
  taosTmrStop(pSyncNode->pElectTimer);
  pSyncNode->pElectTimer = NULL;
1416

M
Minghao Li 已提交
1417 1418 1419 1420 1421 1422 1423 1424 1425 1426
  return ret;
}

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

M
Minghao Li 已提交
1427 1428
int32_t syncNodeResetElectTimer(SSyncNode* pSyncNode) {
  int32_t ret = 0;
M
Minghao Li 已提交
1429 1430 1431 1432 1433 1434 1435
  int32_t electMS;

  if (pSyncNode->pRaftCfg->isStandBy) {
    electMS = TIMER_MAX_MS;
  } else {
    electMS = syncUtilElectRandomMS(pSyncNode->electBaseLine, 2 * pSyncNode->electBaseLine);
  }
M
Minghao Li 已提交
1436
  ret = syncNodeRestartElectTimer(pSyncNode, electMS);
1437 1438 1439 1440 1441 1442 1443 1444

  do {
    char logBuf[128];
    snprintf(logBuf, sizeof(logBuf), "reset elect timer, min:%d, max:%d, ms:%d", pSyncNode->electBaseLine,
             2 * pSyncNode->electBaseLine, electMS);
    syncNodeEventLog(pSyncNode, logBuf);
  } while (0);

M
Minghao Li 已提交
1445 1446 1447
  return ret;
}

M
Minghao Li 已提交
1448
static int32_t syncNodeDoStartHeartbeatTimer(SSyncNode* pSyncNode) {
M
Minghao Li 已提交
1449
  int32_t ret = 0;
S
Shengliang Guan 已提交
1450 1451
  if (syncIsInit()) {
    taosTmrReset(pSyncNode->FpHeartbeatTimerCB, pSyncNode->heartbeatTimerMS, pSyncNode, syncEnv()->pTimerManager,
1452 1453 1454
                 &pSyncNode->pHeartbeatTimer);
    atomic_store_64(&pSyncNode->heartbeatTimerLogicClock, pSyncNode->heartbeatTimerLogicClockUser);
  } else {
M
Minghao Li 已提交
1455
    sError("vgId:%d, start heartbeat timer error, sync env is stop", pSyncNode->vgId);
1456
  }
1457 1458 1459 1460 1461 1462 1463

  do {
    char logBuf[128];
    snprintf(logBuf, sizeof(logBuf), "start heartbeat timer, ms:%d", pSyncNode->heartbeatTimerMS);
    syncNodeEventLog(pSyncNode, logBuf);
  } while (0);

M
Minghao Li 已提交
1464 1465 1466
  return ret;
}

M
Minghao Li 已提交
1467
int32_t syncNodeStartHeartbeatTimer(SSyncNode* pSyncNode) {
1468
  int32_t ret = 0;
M
Minghao Li 已提交
1469

1470
#if 0
M
Minghao Li 已提交
1471
  pSyncNode->heartbeatTimerMS = pSyncNode->hbBaseLine;
1472 1473
  ret = syncNodeDoStartHeartbeatTimer(pSyncNode);
#endif
1474

1475 1476
  for (int i = 0; i < pSyncNode->peersNum; ++i) {
    SSyncTimer* pSyncTimer = syncNodeGetHbTimer(pSyncNode, &(pSyncNode->peersId[i]));
M
Minghao Li 已提交
1477 1478 1479
    if (pSyncTimer != NULL) {
      syncHbTimerStart(pSyncNode, pSyncTimer);
    }
1480
  }
1481

M
Minghao Li 已提交
1482 1483 1484
  return ret;
}

M
Minghao Li 已提交
1485 1486
int32_t syncNodeStopHeartbeatTimer(SSyncNode* pSyncNode) {
  int32_t ret = 0;
1487 1488

#if 0
M
Minghao Li 已提交
1489 1490 1491
  atomic_add_fetch_64(&pSyncNode->heartbeatTimerLogicClockUser, 1);
  taosTmrStop(pSyncNode->pHeartbeatTimer);
  pSyncNode->pHeartbeatTimer = NULL;
1492
#endif
1493

1494 1495
  for (int i = 0; i < pSyncNode->peersNum; ++i) {
    SSyncTimer* pSyncTimer = syncNodeGetHbTimer(pSyncNode, &(pSyncNode->peersId[i]));
M
Minghao Li 已提交
1496 1497 1498
    if (pSyncTimer != NULL) {
      syncHbTimerStop(pSyncNode, pSyncTimer);
    }
1499
  }
1500

M
Minghao Li 已提交
1501 1502 1503
  return ret;
}

1504 1505 1506 1507 1508 1509
int32_t syncNodeRestartHeartbeatTimer(SSyncNode* pSyncNode) {
  syncNodeStopHeartbeatTimer(pSyncNode);
  syncNodeStartHeartbeatTimer(pSyncNode);
  return 0;
}

M
Minghao Li 已提交
1510 1511 1512 1513
// utils --------------
int32_t syncNodeSendMsgById(const SRaftId* destRaftId, SSyncNode* pSyncNode, SRpcMsg* pMsg) {
  SEpSet epSet;
  syncUtilraftId2EpSet(destRaftId, &epSet);
S
Shengliang Guan 已提交
1514
  if (pSyncNode->syncSendMSg != NULL) {
M
Minghao Li 已提交
1515 1516 1517
    // htonl
    syncUtilMsgHtoN(pMsg->pCont);

1518
    pMsg->info.noResp = 1;
S
Shengliang Guan 已提交
1519
    pSyncNode->syncSendMSg(&epSet, pMsg);
M
Minghao Li 已提交
1520
  } else {
M
Minghao Li 已提交
1521 1522
    sError("vgId:%d, sync send msg by id error, fp-send-msg is null", pSyncNode->vgId);
    return -1;
M
Minghao Li 已提交
1523
  }
M
Minghao Li 已提交
1524

M
Minghao Li 已提交
1525 1526 1527 1528 1529 1530
  return 0;
}

int32_t syncNodeSendMsgByInfo(const SNodeInfo* nodeInfo, SSyncNode* pSyncNode, SRpcMsg* pMsg) {
  SEpSet epSet;
  syncUtilnodeInfo2EpSet(nodeInfo, &epSet);
S
Shengliang Guan 已提交
1531
  if (pSyncNode->syncSendMSg != NULL) {
M
Minghao Li 已提交
1532 1533 1534
    // htonl
    syncUtilMsgHtoN(pMsg->pCont);

1535
    pMsg->info.noResp = 1;
S
Shengliang Guan 已提交
1536
    pSyncNode->syncSendMSg(&epSet, pMsg);
M
Minghao Li 已提交
1537
  } else {
M
Minghao Li 已提交
1538
    sError("vgId:%d, sync send msg by info error, fp-send-msg is null", pSyncNode->vgId);
M
Minghao Li 已提交
1539
  }
M
Minghao Li 已提交
1540 1541 1542
  return 0;
}

M
Minghao Li 已提交
1543
cJSON* syncNode2Json(const SSyncNode* pSyncNode) {
C
Cary Xu 已提交
1544
  char   u64buf[128] = {0};
M
Minghao Li 已提交
1545 1546
  cJSON* pRoot = cJSON_CreateObject();

M
Minghao Li 已提交
1547 1548 1549
  if (pSyncNode != NULL) {
    // init by SSyncInfo
    cJSON_AddNumberToObject(pRoot, "vgId", pSyncNode->vgId);
M
Minghao Li 已提交
1550
    cJSON_AddItemToObject(pRoot, "SRaftCfg", raftCfg2Json(pSyncNode->pRaftCfg));
M
Minghao Li 已提交
1551
    cJSON_AddStringToObject(pRoot, "path", pSyncNode->path);
M
Minghao Li 已提交
1552 1553 1554
    cJSON_AddStringToObject(pRoot, "raftStorePath", pSyncNode->raftStorePath);
    cJSON_AddStringToObject(pRoot, "configPath", pSyncNode->configPath);

M
Minghao Li 已提交
1555 1556 1557
    snprintf(u64buf, sizeof(u64buf), "%p", pSyncNode->pWal);
    cJSON_AddStringToObject(pRoot, "pWal", u64buf);

S
Shengliang Guan 已提交
1558
    snprintf(u64buf, sizeof(u64buf), "%p", pSyncNode->msgcb);
M
Minghao Li 已提交
1559
    cJSON_AddStringToObject(pRoot, "rpcClient", u64buf);
S
Shengliang Guan 已提交
1560 1561
    snprintf(u64buf, sizeof(u64buf), "%p", pSyncNode->syncSendMSg);
    cJSON_AddStringToObject(pRoot, "syncSendMSg", u64buf);
M
Minghao Li 已提交
1562

S
Shengliang Guan 已提交
1563
    snprintf(u64buf, sizeof(u64buf), "%p", pSyncNode->msgcb);
M
Minghao Li 已提交
1564
    cJSON_AddStringToObject(pRoot, "queue", u64buf);
S
Shengliang Guan 已提交
1565 1566
    snprintf(u64buf, sizeof(u64buf), "%p", pSyncNode->syncEqMsg);
    cJSON_AddStringToObject(pRoot, "syncEqMsg", u64buf);
M
Minghao Li 已提交
1567 1568 1569 1570 1571 1572 1573 1574 1575 1576 1577 1578 1579 1580 1581 1582 1583 1584

    // init internal
    cJSON* pMe = syncUtilNodeInfo2Json(&pSyncNode->myNodeInfo);
    cJSON_AddItemToObject(pRoot, "myNodeInfo", pMe);
    cJSON* pRaftId = syncUtilRaftId2Json(&pSyncNode->myRaftId);
    cJSON_AddItemToObject(pRoot, "myRaftId", pRaftId);

    cJSON_AddNumberToObject(pRoot, "peersNum", pSyncNode->peersNum);
    cJSON* pPeers = cJSON_CreateArray();
    cJSON_AddItemToObject(pRoot, "peersNodeInfo", pPeers);
    for (int i = 0; i < pSyncNode->peersNum; ++i) {
      cJSON_AddItemToArray(pPeers, syncUtilNodeInfo2Json(&pSyncNode->peersNodeInfo[i]));
    }
    cJSON* pPeersId = cJSON_CreateArray();
    cJSON_AddItemToObject(pRoot, "peersId", pPeersId);
    for (int i = 0; i < pSyncNode->peersNum; ++i) {
      cJSON_AddItemToArray(pPeersId, syncUtilRaftId2Json(&pSyncNode->peersId[i]));
    }
M
Minghao Li 已提交
1585

M
Minghao Li 已提交
1586 1587 1588 1589 1590 1591
    cJSON_AddNumberToObject(pRoot, "replicaNum", pSyncNode->replicaNum);
    cJSON* pReplicasId = cJSON_CreateArray();
    cJSON_AddItemToObject(pRoot, "replicasId", pReplicasId);
    for (int i = 0; i < pSyncNode->replicaNum; ++i) {
      cJSON_AddItemToArray(pReplicasId, syncUtilRaftId2Json(&pSyncNode->replicasId[i]));
    }
M
Minghao Li 已提交
1592

M
Minghao Li 已提交
1593 1594 1595 1596 1597 1598 1599
    // raft algorithm
    snprintf(u64buf, sizeof(u64buf), "%p", pSyncNode->pFsm);
    cJSON_AddStringToObject(pRoot, "pFsm", u64buf);
    cJSON_AddNumberToObject(pRoot, "quorum", pSyncNode->quorum);
    cJSON* pLaderCache = syncUtilRaftId2Json(&pSyncNode->leaderCache);
    cJSON_AddItemToObject(pRoot, "leaderCache", pLaderCache);

M
Minghao Li 已提交
1600
    // life cycle
S
Shengliang Guan 已提交
1601
    snprintf(u64buf, sizeof(u64buf), "%" PRId64, pSyncNode->rid);
M
Minghao Li 已提交
1602 1603
    cJSON_AddStringToObject(pRoot, "rid", u64buf);

M
Minghao Li 已提交
1604 1605 1606
    // tla+ server vars
    cJSON_AddNumberToObject(pRoot, "state", pSyncNode->state);
    cJSON_AddStringToObject(pRoot, "state_str", syncUtilState2String(pSyncNode->state));
M
Minghao Li 已提交
1607
    cJSON_AddItemToObject(pRoot, "pRaftStore", raftStore2Json(pSyncNode->pRaftStore));
M
Minghao Li 已提交
1608 1609 1610 1611 1612 1613 1614 1615 1616 1617 1618

    // tla+ candidate vars
    cJSON_AddItemToObject(pRoot, "pVotesGranted", voteGranted2Json(pSyncNode->pVotesGranted));
    cJSON_AddItemToObject(pRoot, "pVotesRespond", votesRespond2Json(pSyncNode->pVotesRespond));

    // tla+ leader vars
    cJSON_AddItemToObject(pRoot, "pNextIndex", syncIndexMgr2Json(pSyncNode->pNextIndex));
    cJSON_AddItemToObject(pRoot, "pMatchIndex", syncIndexMgr2Json(pSyncNode->pMatchIndex));

    // tla+ log vars
    cJSON_AddItemToObject(pRoot, "pLogStore", logStore2Json(pSyncNode->pLogStore));
S
Shengliang Guan 已提交
1619
    snprintf(u64buf, sizeof(u64buf), "%" PRId64, pSyncNode->commitIndex);
M
Minghao Li 已提交
1620 1621
    cJSON_AddStringToObject(pRoot, "commitIndex", u64buf);

M
Minghao Li 已提交
1622 1623 1624 1625 1626
    // timer ms init
    cJSON_AddNumberToObject(pRoot, "pingBaseLine", pSyncNode->pingBaseLine);
    cJSON_AddNumberToObject(pRoot, "electBaseLine", pSyncNode->electBaseLine);
    cJSON_AddNumberToObject(pRoot, "hbBaseLine", pSyncNode->hbBaseLine);

M
Minghao Li 已提交
1627 1628 1629 1630
    // ping timer
    snprintf(u64buf, sizeof(u64buf), "%p", pSyncNode->pPingTimer);
    cJSON_AddStringToObject(pRoot, "pPingTimer", u64buf);
    cJSON_AddNumberToObject(pRoot, "pingTimerMS", pSyncNode->pingTimerMS);
S
Shengliang Guan 已提交
1631
    snprintf(u64buf, sizeof(u64buf), "%" PRIu64, pSyncNode->pingTimerLogicClock);
M
Minghao Li 已提交
1632
    cJSON_AddStringToObject(pRoot, "pingTimerLogicClock", u64buf);
S
Shengliang Guan 已提交
1633
    snprintf(u64buf, sizeof(u64buf), "%" PRIu64, pSyncNode->pingTimerLogicClockUser);
M
Minghao Li 已提交
1634 1635 1636
    cJSON_AddStringToObject(pRoot, "pingTimerLogicClockUser", u64buf);
    snprintf(u64buf, sizeof(u64buf), "%p", pSyncNode->FpPingTimerCB);
    cJSON_AddStringToObject(pRoot, "FpPingTimerCB", u64buf);
S
Shengliang Guan 已提交
1637
    snprintf(u64buf, sizeof(u64buf), "%" PRIu64, pSyncNode->pingTimerCounter);
M
Minghao Li 已提交
1638 1639 1640 1641 1642 1643
    cJSON_AddStringToObject(pRoot, "pingTimerCounter", u64buf);

    // elect timer
    snprintf(u64buf, sizeof(u64buf), "%p", pSyncNode->pElectTimer);
    cJSON_AddStringToObject(pRoot, "pElectTimer", u64buf);
    cJSON_AddNumberToObject(pRoot, "electTimerMS", pSyncNode->electTimerMS);
S
Shengliang Guan 已提交
1644
    snprintf(u64buf, sizeof(u64buf), "%" PRIu64, pSyncNode->electTimerLogicClock);
M
Minghao Li 已提交
1645 1646 1647
    cJSON_AddStringToObject(pRoot, "electTimerLogicClock", u64buf);
    snprintf(u64buf, sizeof(u64buf), "%p", pSyncNode->FpElectTimerCB);
    cJSON_AddStringToObject(pRoot, "FpElectTimerCB", u64buf);
S
Shengliang Guan 已提交
1648
    snprintf(u64buf, sizeof(u64buf), "%" PRIu64, pSyncNode->electTimerCounter);
M
Minghao Li 已提交
1649 1650 1651 1652 1653 1654
    cJSON_AddStringToObject(pRoot, "electTimerCounter", u64buf);

    // heartbeat timer
    snprintf(u64buf, sizeof(u64buf), "%p", pSyncNode->pHeartbeatTimer);
    cJSON_AddStringToObject(pRoot, "pHeartbeatTimer", u64buf);
    cJSON_AddNumberToObject(pRoot, "heartbeatTimerMS", pSyncNode->heartbeatTimerMS);
S
Shengliang Guan 已提交
1655
    snprintf(u64buf, sizeof(u64buf), "%" PRIu64, pSyncNode->heartbeatTimerLogicClock);
M
Minghao Li 已提交
1656
    cJSON_AddStringToObject(pRoot, "heartbeatTimerLogicClock", u64buf);
S
Shengliang Guan 已提交
1657
    snprintf(u64buf, sizeof(u64buf), "%" PRIu64, pSyncNode->heartbeatTimerLogicClockUser);
M
Minghao Li 已提交
1658 1659 1660
    cJSON_AddStringToObject(pRoot, "heartbeatTimerLogicClockUser", u64buf);
    snprintf(u64buf, sizeof(u64buf), "%p", pSyncNode->FpHeartbeatTimerCB);
    cJSON_AddStringToObject(pRoot, "FpHeartbeatTimerCB", u64buf);
S
Shengliang Guan 已提交
1661
    snprintf(u64buf, sizeof(u64buf), "%" PRIu64, pSyncNode->heartbeatTimerCounter);
M
Minghao Li 已提交
1662 1663 1664 1665 1666 1667 1668 1669 1670 1671 1672 1673 1674 1675 1676 1677 1678
    cJSON_AddStringToObject(pRoot, "heartbeatTimerCounter", u64buf);

    // callback
    snprintf(u64buf, sizeof(u64buf), "%p", pSyncNode->FpOnPing);
    cJSON_AddStringToObject(pRoot, "FpOnPing", u64buf);
    snprintf(u64buf, sizeof(u64buf), "%p", pSyncNode->FpOnPingReply);
    cJSON_AddStringToObject(pRoot, "FpOnPingReply", u64buf);
    snprintf(u64buf, sizeof(u64buf), "%p", pSyncNode->FpOnRequestVote);
    cJSON_AddStringToObject(pRoot, "FpOnRequestVote", u64buf);
    snprintf(u64buf, sizeof(u64buf), "%p", pSyncNode->FpOnRequestVoteReply);
    cJSON_AddStringToObject(pRoot, "FpOnRequestVoteReply", u64buf);
    snprintf(u64buf, sizeof(u64buf), "%p", pSyncNode->FpOnAppendEntries);
    cJSON_AddStringToObject(pRoot, "FpOnAppendEntries", u64buf);
    snprintf(u64buf, sizeof(u64buf), "%p", pSyncNode->FpOnAppendEntriesReply);
    cJSON_AddStringToObject(pRoot, "FpOnAppendEntriesReply", u64buf);
    snprintf(u64buf, sizeof(u64buf), "%p", pSyncNode->FpOnTimeout);
    cJSON_AddStringToObject(pRoot, "FpOnTimeout", u64buf);
M
Minghao Li 已提交
1679 1680 1681 1682 1683 1684 1685 1686 1687 1688 1689 1690 1691

    // restoreFinish
    cJSON_AddNumberToObject(pRoot, "restoreFinish", pSyncNode->restoreFinish);

    // snapshot senders
    cJSON* pSenders = cJSON_CreateArray();
    cJSON_AddItemToObject(pRoot, "senders", pSenders);
    for (int i = 0; i < TSDB_MAX_REPLICA; ++i) {
      cJSON_AddItemToArray(pSenders, snapshotSender2Json((pSyncNode->senders)[i]));
    }

    // snapshot receivers
    cJSON* pReceivers = cJSON_CreateArray();
1692
    cJSON_AddItemToObject(pRoot, "receiver", snapshotReceiver2Json(pSyncNode->pNewNodeReceiver));
M
Minghao Li 已提交
1693 1694 1695

    // changing
    cJSON_AddNumberToObject(pRoot, "changing", pSyncNode->changing);
M
Minghao Li 已提交
1696 1697 1698 1699 1700 1701 1702
  }

  cJSON* pJson = cJSON_CreateObject();
  cJSON_AddItemToObject(pJson, "SSyncNode", pRoot);
  return pJson;
}

M
Minghao Li 已提交
1703 1704 1705 1706 1707 1708 1709
char* syncNode2Str(const SSyncNode* pSyncNode) {
  cJSON* pJson = syncNode2Json(pSyncNode);
  char*  serialized = cJSON_Print(pJson);
  cJSON_Delete(pJson);
  return serialized;
}

1710
inline void syncNodeEventLog(const SSyncNode* pSyncNode, char* str) {
M
Minghao Li 已提交
1711 1712 1713 1714
  if (pSyncNode == NULL) {
    return;
  }

M
Minghao Li 已提交
1715
  SSnapshot snapshot = {.data = NULL, .lastApplyIndex = -1, .lastApplyTerm = 0};
M
Minghao Li 已提交
1716
  if (pSyncNode->pFsm != NULL && pSyncNode->pFsm->FpGetSnapshotInfo != NULL) {
M
Minghao Li 已提交
1717 1718
    pSyncNode->pFsm->FpGetSnapshotInfo(pSyncNode->pFsm, &snapshot);
  }
M
Minghao Li 已提交
1719 1720 1721 1722 1723 1724 1725

  SyncIndex logLastIndex = SYNC_INDEX_INVALID;
  SyncIndex logBeginIndex = SYNC_INDEX_INVALID;
  if (pSyncNode->pLogStore != NULL) {
    logLastIndex = pSyncNode->pLogStore->syncLogLastIndex(pSyncNode->pLogStore);
    logBeginIndex = pSyncNode->pLogStore->syncLogBeginIndex(pSyncNode->pLogStore);
  }
M
Minghao Li 已提交
1726

M
Minghao Li 已提交
1727
  char* pCfgStr = syncCfg2SimpleStr(&(pSyncNode->pRaftCfg->cfg));
1728 1729 1730 1731
  char* printStr = "";
  if (pCfgStr != NULL) {
    printStr = pCfgStr;
  }
M
Minghao Li 已提交
1732

1733 1734 1735
  char*   peerStateStr = syncNodePeerState2Str(pSyncNode);
  int32_t userStrLen = strlen(str) + strlen(peerStateStr);

M
Minghao Li 已提交
1736
  if (userStrLen < 256) {
M
Minghao Li 已提交
1737
    char logBuf[256 + 256];
1738 1739
    if (pSyncNode != NULL && pSyncNode->pRaftCfg != NULL && pSyncNode->pRaftStore != NULL) {
      snprintf(logBuf, sizeof(logBuf),
M
Minghao Li 已提交
1740 1741
               "vgId:%d, sync %s %s, tm:%" PRIu64 ", cmt:%" PRId64 ", fst:%" PRId64 ", lst:%" PRId64 ", min:%" PRId64
               ", snap:%" PRId64 ", snap-tm:%" PRIu64
M
Minghao Li 已提交
1742 1743 1744
               ", sby:%d, "
               "stgy:%d, bch:%d, "
               "r-num:%d, "
1745
               "lcfg:%" PRId64 ", chging:%d, rsto:%d, dquorum:%d, elt:%" PRId64 ", hb:%" PRId64 ", %s, %s",
1746
               pSyncNode->vgId, syncUtilState2String(pSyncNode->state), str, pSyncNode->pRaftStore->currentTerm,
M
Minghao Li 已提交
1747 1748 1749 1750
               pSyncNode->commitIndex, logBeginIndex, logLastIndex, pSyncNode->minMatchIndex, snapshot.lastApplyIndex,
               snapshot.lastApplyTerm, pSyncNode->pRaftCfg->isStandBy, pSyncNode->pRaftCfg->snapshotStrategy,
               pSyncNode->pRaftCfg->batchSize, pSyncNode->replicaNum, pSyncNode->pRaftCfg->lastConfigIndex,
               pSyncNode->changing, pSyncNode->restoreFinish, syncNodeDynamicQuorum(pSyncNode),
M
Minghao Li 已提交
1751
               pSyncNode->electTimerLogicClock, pSyncNode->heartbeatTimerLogicClockUser, peerStateStr, printStr);
1752 1753 1754
    } else {
      snprintf(logBuf, sizeof(logBuf), "%s", str);
    }
1755
    // sDebug("%s", logBuf);
M
Minghao Li 已提交
1756 1757
    // sInfo("%s", logBuf);
    sTrace("%s", logBuf);
M
Minghao Li 已提交
1758

M
Minghao Li 已提交
1759
  } else {
M
Minghao Li 已提交
1760
    int   len = 256 + userStrLen;
M
Minghao Li 已提交
1761
    char* s = (char*)taosMemoryMalloc(len);
1762 1763
    if (pSyncNode != NULL && pSyncNode->pRaftCfg != NULL && pSyncNode->pRaftStore != NULL) {
      snprintf(s, len,
M
Minghao Li 已提交
1764 1765
               "vgId:%d, sync %s %s, tm:%" PRIu64 ", cmt:%" PRId64 ", fst:%" PRId64 ", lst:%" PRId64 ", min:%" PRId64
               ", snap:%" PRId64 ", snap-tm:%" PRIu64
M
Minghao Li 已提交
1766 1767 1768
               ", sby:%d, "
               "stgy:%d, bch:%d, "
               "r-num:%d, "
1769
               "lcfg:%" PRId64 ", chging:%d, rsto:%d, dquorum:%d, elt:%" PRId64 ", hb:%" PRId64 ", %s, %s",
1770
               pSyncNode->vgId, syncUtilState2String(pSyncNode->state), str, pSyncNode->pRaftStore->currentTerm,
M
Minghao Li 已提交
1771 1772 1773 1774
               pSyncNode->commitIndex, logBeginIndex, logLastIndex, pSyncNode->minMatchIndex, snapshot.lastApplyIndex,
               snapshot.lastApplyTerm, pSyncNode->pRaftCfg->isStandBy, pSyncNode->pRaftCfg->snapshotStrategy,
               pSyncNode->pRaftCfg->batchSize, pSyncNode->replicaNum, pSyncNode->pRaftCfg->lastConfigIndex,
               pSyncNode->changing, pSyncNode->restoreFinish, syncNodeDynamicQuorum(pSyncNode),
M
Minghao Li 已提交
1775
               pSyncNode->electTimerLogicClock, pSyncNode->heartbeatTimerLogicClockUser, peerStateStr, printStr);
1776 1777 1778
    } else {
      snprintf(s, len, "%s", str);
    }
1779
    // sDebug("%s", s);
M
Minghao Li 已提交
1780 1781
    // sInfo("%s", s);
    sTrace("%s", s);
M
Minghao Li 已提交
1782 1783
    taosMemoryFree(s);
  }
M
Minghao Li 已提交
1784

M
Minghao Li 已提交
1785
  taosMemoryFree(peerStateStr);
M
Minghao Li 已提交
1786
  taosMemoryFree(pCfgStr);
M
Minghao Li 已提交
1787 1788
}

1789
inline void syncNodeErrorLog(const SSyncNode* pSyncNode, char* str) {
M
Minghao Li 已提交
1790 1791 1792 1793
  if (pSyncNode == NULL) {
    return;
  }

M
Minghao Li 已提交
1794 1795 1796
  int32_t userStrLen = strlen(str);

  SSnapshot snapshot = {.data = NULL, .lastApplyIndex = -1, .lastApplyTerm = 0};
M
Minghao Li 已提交
1797
  if (pSyncNode->pFsm != NULL && pSyncNode->pFsm->FpGetSnapshotInfo != NULL) {
M
Minghao Li 已提交
1798 1799
    pSyncNode->pFsm->FpGetSnapshotInfo(pSyncNode->pFsm, &snapshot);
  }
M
Minghao Li 已提交
1800 1801 1802 1803 1804 1805 1806 1807 1808 1809 1810 1811 1812

  SyncIndex logLastIndex = SYNC_INDEX_INVALID;
  SyncIndex logBeginIndex = SYNC_INDEX_INVALID;
  if (pSyncNode->pLogStore != NULL) {
    logLastIndex = pSyncNode->pLogStore->syncLogLastIndex(pSyncNode->pLogStore);
    logBeginIndex = pSyncNode->pLogStore->syncLogBeginIndex(pSyncNode->pLogStore);
  }

  char* pCfgStr = syncCfg2SimpleStr(&(pSyncNode->pRaftCfg->cfg));
  char* printStr = "";
  if (pCfgStr != NULL) {
    printStr = pCfgStr;
  }
M
Minghao Li 已提交
1813 1814

  if (userStrLen < 256) {
M
Minghao Li 已提交
1815
    char logBuf[256 + 256];
1816 1817
    if (pSyncNode != NULL && pSyncNode->pRaftCfg != NULL && pSyncNode->pRaftStore != NULL) {
      snprintf(logBuf, sizeof(logBuf),
M
Minghao Li 已提交
1818 1819
               "vgId:%d, sync %s %s, tm:%" PRIu64 ", cmt:%" PRId64 ", fst:%" PRId64 ", lst:%" PRId64 ", min:%" PRId64
               ", snap:%" PRId64 ", snap-tm:%" PRIu64
M
Minghao Li 已提交
1820 1821 1822
               ", sby:%d, "
               "stgy:%d, bch:%d, "
               "r-num:%d, "
M
Minghao Li 已提交
1823
               "lcfg:%" PRId64 ", chging:%d, rsto:%d, dquorum:%d, elt:%" PRId64 ", hb:%" PRId64 ", %s",
1824
               pSyncNode->vgId, syncUtilState2String(pSyncNode->state), str, pSyncNode->pRaftStore->currentTerm,
M
Minghao Li 已提交
1825 1826 1827 1828
               pSyncNode->commitIndex, logBeginIndex, logLastIndex, pSyncNode->minMatchIndex, snapshot.lastApplyIndex,
               snapshot.lastApplyTerm, pSyncNode->pRaftCfg->isStandBy, pSyncNode->pRaftCfg->snapshotStrategy,
               pSyncNode->pRaftCfg->batchSize, pSyncNode->replicaNum, pSyncNode->pRaftCfg->lastConfigIndex,
               pSyncNode->changing, pSyncNode->restoreFinish, syncNodeDynamicQuorum(pSyncNode),
M
Minghao Li 已提交
1829
               pSyncNode->electTimerLogicClock, pSyncNode->heartbeatTimerLogicClockUser, printStr);
1830 1831 1832
    } else {
      snprintf(logBuf, sizeof(logBuf), "%s", str);
    }
M
Minghao Li 已提交
1833 1834 1835
    sError("%s", logBuf);

  } else {
M
Minghao Li 已提交
1836
    int   len = 256 + userStrLen;
M
Minghao Li 已提交
1837
    char* s = (char*)taosMemoryMalloc(len);
1838 1839
    if (pSyncNode != NULL && pSyncNode->pRaftCfg != NULL && pSyncNode->pRaftStore != NULL) {
      snprintf(s, len,
M
Minghao Li 已提交
1840 1841
               "vgId:%d, sync %s %s, tm:%" PRIu64 ", cmt:%" PRId64 ", fst:%" PRId64 ", lst:%" PRId64 ", min:%" PRId64
               ", snap:%" PRId64 ", snap-tm:%" PRIu64
M
Minghao Li 已提交
1842 1843 1844
               ", sby:%d, "
               "stgy:%d, bch:%d, "
               "r-num:%d, "
M
Minghao Li 已提交
1845
               "lcfg:%" PRId64 ", chging:%d, rsto:%d, dquorum:%d, elt:%" PRId64 ", hb:%" PRId64 ", %s",
1846
               pSyncNode->vgId, syncUtilState2String(pSyncNode->state), str, pSyncNode->pRaftStore->currentTerm,
M
Minghao Li 已提交
1847 1848 1849 1850
               pSyncNode->commitIndex, logBeginIndex, logLastIndex, pSyncNode->minMatchIndex, snapshot.lastApplyIndex,
               snapshot.lastApplyTerm, pSyncNode->pRaftCfg->isStandBy, pSyncNode->pRaftCfg->snapshotStrategy,
               pSyncNode->pRaftCfg->batchSize, pSyncNode->replicaNum, pSyncNode->pRaftCfg->lastConfigIndex,
               pSyncNode->changing, pSyncNode->restoreFinish, syncNodeDynamicQuorum(pSyncNode),
M
Minghao Li 已提交
1851
               pSyncNode->electTimerLogicClock, pSyncNode->heartbeatTimerLogicClockUser, printStr);
1852 1853 1854
    } else {
      snprintf(s, len, "%s", str);
    }
M
Minghao Li 已提交
1855 1856 1857
    sError("%s", s);
    taosMemoryFree(s);
  }
M
Minghao Li 已提交
1858 1859

  taosMemoryFree(pCfgStr);
M
Minghao Li 已提交
1860 1861
}

1862
inline char* syncNode2SimpleStr(const SSyncNode* pSyncNode) {
M
Minghao Li 已提交
1863 1864
  int   len = 256;
  char* s = (char*)taosMemoryMalloc(len);
M
Minghao Li 已提交
1865 1866 1867 1868 1869 1870 1871 1872

  SSnapshot snapshot = {.data = NULL, .lastApplyIndex = -1, .lastApplyTerm = 0};
  if (pSyncNode->pFsm->FpGetSnapshotInfo != NULL) {
    pSyncNode->pFsm->FpGetSnapshotInfo(pSyncNode->pFsm, &snapshot);
  }
  SyncIndex logLastIndex = pSyncNode->pLogStore->syncLogLastIndex(pSyncNode->pLogStore);
  SyncIndex logBeginIndex = pSyncNode->pLogStore->syncLogBeginIndex(pSyncNode->pLogStore);

M
Minghao Li 已提交
1873
  snprintf(s, len,
M
Minghao Li 已提交
1874 1875 1876 1877
           "vgId:%d, sync %s, tm:%" PRIu64 ", cmt:%" PRId64 ", fst:%" PRId64 ", lst:%" PRId64 ", snap:%" PRId64
           ", sby:%d, "
           "r-num:%d, "
           "lcfg:%" PRId64 ", chging:%d, rsto:%d",
M
Minghao Li 已提交
1878 1879 1880 1881
           pSyncNode->vgId, syncUtilState2String(pSyncNode->state), pSyncNode->pRaftStore->currentTerm,
           pSyncNode->commitIndex, logBeginIndex, logLastIndex, snapshot.lastApplyIndex, pSyncNode->pRaftCfg->isStandBy,
           pSyncNode->replicaNum, pSyncNode->pRaftCfg->lastConfigIndex, pSyncNode->changing, pSyncNode->restoreFinish);

M
Minghao Li 已提交
1882 1883 1884
  return s;
}

1885
inline bool syncNodeInConfig(SSyncNode* pSyncNode, const SSyncCfg* config) {
1886 1887 1888 1889 1890 1891 1892 1893 1894 1895 1896 1897 1898 1899 1900 1901 1902 1903 1904 1905 1906 1907 1908 1909 1910 1911
  bool b1 = false;
  bool b2 = false;

  for (int i = 0; i < config->replicaNum; ++i) {
    if (strcmp((config->nodeInfo)[i].nodeFqdn, pSyncNode->myNodeInfo.nodeFqdn) == 0 &&
        (config->nodeInfo)[i].nodePort == pSyncNode->myNodeInfo.nodePort) {
      b1 = true;
      break;
    }
  }

  for (int i = 0; i < config->replicaNum; ++i) {
    SRaftId raftId;
    raftId.addr = syncUtilAddr2U64((config->nodeInfo)[i].nodeFqdn, (config->nodeInfo)[i].nodePort);
    raftId.vgId = pSyncNode->vgId;

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

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

1912 1913 1914 1915 1916 1917 1918 1919 1920 1921 1922 1923 1924
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 已提交
1925
void syncNodeDoConfigChange(SSyncNode* pSyncNode, SSyncCfg* pNewConfig, SyncIndex lastConfigChangeIndex) {
1926
  SSyncCfg oldConfig = pSyncNode->pRaftCfg->cfg;
1927 1928 1929 1930
  if (!syncIsConfigChanged(&oldConfig, pNewConfig)) {
    sInfo("vgId:1, sync not reconfig since not changed");
    return;
  }
S
Shengliang Guan 已提交
1931

1932
  pSyncNode->pRaftCfg->cfg = *pNewConfig;
1933 1934
  pSyncNode->pRaftCfg->lastConfigIndex = lastConfigChangeIndex;

M
Minghao Li 已提交
1935 1936
  bool IamInOld = syncNodeInConfig(pSyncNode, &oldConfig);
  bool IamInNew = syncNodeInConfig(pSyncNode, pNewConfig);
M
Minghao Li 已提交
1937

M
Minghao Li 已提交
1938 1939
  bool isDrop = false;
  bool isAdd = false;
M
Minghao Li 已提交
1940

M
Minghao Li 已提交
1941 1942 1943 1944
  if (IamInOld && !IamInNew) {
    isDrop = true;
  } else {
    isDrop = false;
1945
  }
1946

M
Minghao Li 已提交
1947 1948 1949 1950 1951
  if (!IamInOld && IamInNew) {
    isAdd = true;
  } else {
    isAdd = false;
  }
M
Minghao Li 已提交
1952

M
Minghao Li 已提交
1953 1954 1955 1956 1957 1958 1959 1960 1961 1962 1963
  // log begin config change
  do {
    char  eventLog[256];
    char* pOldCfgStr = syncCfg2SimpleStr(&oldConfig);
    char* pNewCfgStr = syncCfg2SimpleStr(pNewConfig);
    snprintf(eventLog, sizeof(eventLog), "begin do config change, from %s to %s", pOldCfgStr, pNewCfgStr);
    syncNodeEventLog(pSyncNode, eventLog);
    taosMemoryFree(pOldCfgStr);
    taosMemoryFree(pNewCfgStr);
  } while (0);

M
Minghao Li 已提交
1964 1965
  if (IamInNew) {
    pSyncNode->pRaftCfg->isStandBy = 0;  // change isStandBy to normal
M
Minghao Li 已提交
1966
  }
M
Minghao Li 已提交
1967 1968
  if (isDrop) {
    pSyncNode->pRaftCfg->isStandBy = 1;  // set standby
M
Minghao Li 已提交
1969 1970
  }

M
Minghao Li 已提交
1971
  // add last config index
M
Minghao Li 已提交
1972
  raftCfgAddConfigIndex(pSyncNode->pRaftCfg, lastConfigChangeIndex);
M
Minghao Li 已提交
1973

M
Minghao Li 已提交
1974 1975 1976 1977 1978 1979 1980 1981 1982 1983 1984
  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];
    for (int i = 0; i < TSDB_MAX_REPLICA; ++i) {
      oldSenders[i] = (pSyncNode->senders)[i];
M
Minghao Li 已提交
1985

M
Minghao Li 已提交
1986 1987 1988 1989
      char* eventLog = snapshotSender2SimpleStr(oldSenders[i], "snapshot sender save old");
      syncNodeEventLog(pSyncNode, eventLog);
      taosMemoryFree(eventLog);
    }
1990

M
Minghao Li 已提交
1991 1992 1993 1994 1995 1996 1997 1998 1999 2000 2001 2002 2003 2004 2005 2006
    // init internal
    pSyncNode->myNodeInfo = pSyncNode->pRaftCfg->cfg.nodeInfo[pSyncNode->pRaftCfg->cfg.myIndex];
    syncUtilnodeInfo2raftId(&pSyncNode->myNodeInfo, pSyncNode->vgId, &pSyncNode->myRaftId);

    // init peersNum, peers, peersId
    pSyncNode->peersNum = pSyncNode->pRaftCfg->cfg.replicaNum - 1;
    int j = 0;
    for (int i = 0; i < pSyncNode->pRaftCfg->cfg.replicaNum; ++i) {
      if (i != pSyncNode->pRaftCfg->cfg.myIndex) {
        pSyncNode->peersNodeInfo[j] = pSyncNode->pRaftCfg->cfg.nodeInfo[i];
        j++;
      }
    }
    for (int i = 0; i < pSyncNode->peersNum; ++i) {
      syncUtilnodeInfo2raftId(&pSyncNode->peersNodeInfo[i], pSyncNode->vgId, &pSyncNode->peersId[i]);
    }
2007

M
Minghao Li 已提交
2008 2009 2010 2011 2012
    // init replicaNum, replicasId
    pSyncNode->replicaNum = pSyncNode->pRaftCfg->cfg.replicaNum;
    for (int i = 0; i < pSyncNode->pRaftCfg->cfg.replicaNum; ++i) {
      syncUtilnodeInfo2raftId(&pSyncNode->pRaftCfg->cfg.nodeInfo[i], pSyncNode->vgId, &pSyncNode->replicasId[i]);
    }
2013

2014 2015 2016
    // update quorum first
    pSyncNode->quorum = syncUtilQuorum(pSyncNode->pRaftCfg->cfg.replicaNum);

M
Minghao Li 已提交
2017 2018 2019 2020
    syncIndexMgrUpdate(pSyncNode->pNextIndex, pSyncNode);
    syncIndexMgrUpdate(pSyncNode->pMatchIndex, pSyncNode);
    voteGrantedUpdate(pSyncNode->pVotesGranted, pSyncNode);
    votesRespondUpdate(pSyncNode->pVotesRespond, pSyncNode);
M
Minghao Li 已提交
2021

M
Minghao Li 已提交
2022
    // reset snapshot senders
2023

M
Minghao Li 已提交
2024 2025 2026 2027
    // clear new
    for (int i = 0; i < TSDB_MAX_REPLICA; ++i) {
      (pSyncNode->senders)[i] = NULL;
    }
M
Minghao Li 已提交
2028

M
Minghao Li 已提交
2029 2030 2031 2032 2033 2034 2035 2036 2037 2038 2039 2040
    // reset new
    for (int i = 0; i < pSyncNode->replicaNum; ++i) {
      // reset sender
      bool reset = false;
      for (int j = 0; j < TSDB_MAX_REPLICA; ++j) {
        if (syncUtilSameId(&(pSyncNode->replicasId)[i], &oldReplicasId[j])) {
          char     host[128];
          uint16_t port;
          syncUtilU642Addr((pSyncNode->replicasId)[i].addr, host, sizeof(host), &port);

          do {
            char eventLog[256];
S
Shengliang Guan 已提交
2041
            snprintf(eventLog, sizeof(eventLog), "snapshot sender reset for: %" PRIu64 ", newIndex:%d, %s:%d, %p",
M
Minghao Li 已提交
2042 2043 2044 2045 2046 2047 2048 2049 2050 2051 2052 2053 2054 2055 2056 2057 2058 2059 2060 2061
                     (pSyncNode->replicasId)[i].addr, i, host, port, oldSenders[j]);
            syncNodeEventLog(pSyncNode, eventLog);
          } while (0);

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

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

          do {
            char eventLog[256];
            snprintf(eventLog, sizeof(eventLog),
                     "snapshot sender udpate replicaIndex from %d to %d, %s:%d, %p, reset:%d", oldreplicaIndex, i, host,
                     port, (pSyncNode->senders)[i], reset);
            syncNodeEventLog(pSyncNode, eventLog);
          } while (0);
        }
2062 2063
      }
    }
2064

M
Minghao Li 已提交
2065 2066 2067 2068
    // create new
    for (int i = 0; i < TSDB_MAX_REPLICA; ++i) {
      if ((pSyncNode->senders)[i] == NULL) {
        (pSyncNode->senders)[i] = snapshotSenderCreate(pSyncNode, i);
M
Minghao Li 已提交
2069

M
Minghao Li 已提交
2070 2071 2072 2073
        char* eventLog = snapshotSender2SimpleStr((pSyncNode->senders)[i], "snapshot sender create new");
        syncNodeEventLog(pSyncNode, eventLog);
        taosMemoryFree(eventLog);
      }
2074 2075
    }

M
Minghao Li 已提交
2076 2077 2078 2079
    // free old
    for (int i = 0; i < TSDB_MAX_REPLICA; ++i) {
      if (oldSenders[i] != NULL) {
        snapshotSenderDestroy(oldSenders[i]);
M
Minghao Li 已提交
2080

M
Minghao Li 已提交
2081 2082 2083 2084 2085
        do {
          char eventLog[128];
          snprintf(eventLog, sizeof(eventLog), "snapshot sender delete old %p replica-index:%d", oldSenders[i], i);
          syncNodeEventLog(pSyncNode, eventLog);
        } while (0);
M
Minghao Li 已提交
2086

M
Minghao Li 已提交
2087 2088
        oldSenders[i] = NULL;
      }
2089 2090
    }

2091
    // persist cfg
M
Minghao Li 已提交
2092
    raftCfgPersist(pSyncNode->pRaftCfg);
2093

M
Minghao Li 已提交
2094 2095 2096
    char  tmpbuf[512];
    char* oldStr = syncCfg2SimpleStr(&oldConfig);
    char* newStr = syncCfg2SimpleStr(pNewConfig);
2097 2098
    snprintf(tmpbuf, sizeof(tmpbuf), "config change from %d to %d, index:%" PRId64 ", %s  -->  %s",
             oldConfig.replicaNum, pNewConfig->replicaNum, lastConfigChangeIndex, oldStr, newStr);
M
Minghao Li 已提交
2099 2100
    taosMemoryFree(oldStr);
    taosMemoryFree(newStr);
M
Minghao Li 已提交
2101

M
Minghao Li 已提交
2102 2103 2104
    // change isStandBy to normal (election timeout)
    if (pSyncNode->state == TAOS_SYNC_STATE_LEADER) {
      syncNodeBecomeLeader(pSyncNode, tmpbuf);
2105 2106 2107 2108 2109

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

M
Minghao Li 已提交
2110 2111 2112 2113
    } else {
      syncNodeBecomeFollower(pSyncNode, tmpbuf);
    }
  } else {
2114
    // persist cfg
M
Minghao Li 已提交
2115
    raftCfgPersist(pSyncNode->pRaftCfg);
2116

M
Minghao Li 已提交
2117 2118 2119
    char  tmpbuf[512];
    char* oldStr = syncCfg2SimpleStr(&oldConfig);
    char* newStr = syncCfg2SimpleStr(pNewConfig);
2120 2121
    snprintf(tmpbuf, sizeof(tmpbuf), "do not config change from %d to %d, index:%" PRId64 ", %s  -->  %s",
             oldConfig.replicaNum, pNewConfig->replicaNum, lastConfigChangeIndex, oldStr, newStr);
M
Minghao Li 已提交
2122 2123 2124
    taosMemoryFree(oldStr);
    taosMemoryFree(newStr);
    syncNodeEventLog(pSyncNode, tmpbuf);
2125
  }
2126

M
Minghao Li 已提交
2127
_END:
M
Minghao Li 已提交
2128 2129 2130 2131 2132 2133 2134 2135 2136 2137 2138

  // log end config change
  do {
    char  eventLog[256];
    char* pOldCfgStr = syncCfg2SimpleStr(&oldConfig);
    char* pNewCfgStr = syncCfg2SimpleStr(pNewConfig);
    snprintf(eventLog, sizeof(eventLog), "end do config change, from %s to %s", pOldCfgStr, pNewCfgStr);
    syncNodeEventLog(pSyncNode, eventLog);
    taosMemoryFree(pOldCfgStr);
    taosMemoryFree(pNewCfgStr);
  } while (0);
M
Minghao Li 已提交
2139
  return;
M
Minghao Li 已提交
2140 2141
}

M
Minghao Li 已提交
2142 2143 2144 2145
// raft state change --------------
void syncNodeUpdateTerm(SSyncNode* pSyncNode, SyncTerm term) {
  if (term > pSyncNode->pRaftStore->currentTerm) {
    raftStoreSetTerm(pSyncNode->pRaftStore, term);
2146
    char tmpBuf[64];
S
Shengliang Guan 已提交
2147
    snprintf(tmpBuf, sizeof(tmpBuf), "update term to %" PRIu64, term);
2148
    syncNodeBecomeFollower(pSyncNode, tmpBuf);
M
Minghao Li 已提交
2149 2150 2151 2152
    raftStoreClearVote(pSyncNode->pRaftStore);
  }
}

2153 2154 2155 2156 2157 2158
void syncNodeUpdateTermWithoutStepDown(SSyncNode* pSyncNode, SyncTerm term) {
  if (term > pSyncNode->pRaftStore->currentTerm) {
    raftStoreSetTerm(pSyncNode->pRaftStore, term);
  }
}

M
Minghao Li 已提交
2159 2160 2161 2162 2163
void syncNodeStepDown(SSyncNode* pSyncNode, SyncTerm newTerm) {
  ASSERT(pSyncNode->pRaftStore->currentTerm <= newTerm);

  do {
    char logBuf[128];
S
Shengliang Guan 已提交
2164
    snprintf(logBuf, sizeof(logBuf), "step down, new-term:%" PRIu64 ", current-term:%" PRIu64, newTerm,
M
Minghao Li 已提交
2165 2166 2167 2168 2169 2170 2171 2172 2173 2174 2175 2176 2177 2178 2179 2180 2181 2182
             pSyncNode->pRaftStore->currentTerm);
    syncNodeEventLog(pSyncNode, logBuf);
  } while (0);

  if (pSyncNode->pRaftStore->currentTerm < newTerm) {
    raftStoreSetTerm(pSyncNode->pRaftStore, newTerm);
    char tmpBuf[64];
    snprintf(tmpBuf, sizeof(tmpBuf), "step down, update term to %" PRIu64, newTerm);
    syncNodeBecomeFollower(pSyncNode, tmpBuf);
    raftStoreClearVote(pSyncNode->pRaftStore);

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

2183 2184
void syncNodeLeaderChangeRsp(SSyncNode* pSyncNode) { syncRespCleanRsp(pSyncNode->pSyncRespMgr); }

2185
void syncNodeBecomeFollower(SSyncNode* pSyncNode, const char* debugStr) {
M
Minghao Li 已提交
2186
  // maybe clear leader cache
M
Minghao Li 已提交
2187 2188 2189 2190
  if (pSyncNode->state == TAOS_SYNC_STATE_LEADER) {
    pSyncNode->leaderCache = EMPTY_RAFT_ID;
  }

M
Minghao Li 已提交
2191
  // state change
M
Minghao Li 已提交
2192 2193 2194
  pSyncNode->state = TAOS_SYNC_STATE_FOLLOWER;
  syncNodeStopHeartbeatTimer(pSyncNode);

M
Minghao Li 已提交
2195 2196
  // reset elect timer
  syncNodeResetElectTimer(pSyncNode);
M
Minghao Li 已提交
2197

2198 2199 2200
  // send rsp to client
  syncNodeLeaderChangeRsp(pSyncNode);

2201 2202 2203 2204 2205
  // call back
  if (pSyncNode->pFsm != NULL && pSyncNode->pFsm->FpBecomeFollowerCb != NULL) {
    pSyncNode->pFsm->FpBecomeFollowerCb(pSyncNode->pFsm);
  }

M
Minghao Li 已提交
2206 2207 2208
  // min match index
  pSyncNode->minMatchIndex = SYNC_INDEX_INVALID;

M
Minghao Li 已提交
2209 2210 2211 2212 2213 2214 2215 2216 2217 2218 2219 2220 2221 2222
  // trace log
  do {
    int32_t debugStrLen = strlen(debugStr);
    if (debugStrLen < 256) {
      char eventLog[256 + 64];
      snprintf(eventLog, sizeof(eventLog), "become follower %s", debugStr);
      syncNodeEventLog(pSyncNode, eventLog);
    } else {
      char* eventLog = taosMemoryMalloc(debugStrLen + 64);
      snprintf(eventLog, debugStrLen, "become follower %s", debugStr);
      syncNodeEventLog(pSyncNode, eventLog);
      taosMemoryFree(eventLog);
    }
  } while (0);
M
Minghao Li 已提交
2223 2224 2225 2226 2227 2228 2229 2230 2231 2232 2233 2234 2235 2236 2237 2238 2239 2240 2241 2242
}

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

2246 2247 2248
  // reset restoreFinish
  pSyncNode->restoreFinish = false;

M
Minghao Li 已提交
2249
  // state change
M
Minghao Li 已提交
2250
  pSyncNode->state = TAOS_SYNC_STATE_LEADER;
M
Minghao Li 已提交
2251 2252

  // set leader cache
M
Minghao Li 已提交
2253 2254 2255
  pSyncNode->leaderCache = pSyncNode->myRaftId;

  for (int i = 0; i < pSyncNode->pNextIndex->replicaNum; ++i) {
M
Minghao Li 已提交
2256 2257
    // maybe overwrite myself, no harm
    // just do it!
2258 2259 2260 2261 2262 2263 2264 2265 2266

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

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

  for (int i = 0; i < pSyncNode->pMatchIndex->replicaNum; ++i) {
M
Minghao Li 已提交
2270 2271
    // maybe overwrite myself, no harm
    // just do it!
M
Minghao Li 已提交
2272 2273 2274
    pSyncNode->pMatchIndex->index[i] = SYNC_INDEX_INVALID;
  }

M
Minghao Li 已提交
2275 2276 2277
  // init peer mgr
  syncNodePeerStateInit(pSyncNode);

2278 2279
  // update sender private term
  SSyncSnapshotSender* pMySender = syncNodeGetSnapshotSender(pSyncNode, &(pSyncNode->myRaftId));
2280 2281 2282 2283 2284
  if (pMySender != NULL) {
    for (int i = 0; i < pSyncNode->pMatchIndex->replicaNum; ++i) {
      if ((pSyncNode->senders)[i]->privateTerm > pMySender->privateTerm) {
        pMySender->privateTerm = (pSyncNode->senders)[i]->privateTerm;
      }
2285
    }
2286
    (pMySender->privateTerm) += 100;
2287 2288
  }

2289 2290 2291 2292 2293
  // close receiver
  if (snapshotReceiverIsStart(pSyncNode->pNewNodeReceiver)) {
    snapshotReceiverForceStop(pSyncNode->pNewNodeReceiver);
  }

M
Minghao Li 已提交
2294
  // stop elect timer
M
Minghao Li 已提交
2295
  syncNodeStopElectTimer(pSyncNode);
M
Minghao Li 已提交
2296

M
Minghao Li 已提交
2297 2298
  // start heartbeat timer
  syncNodeStartHeartbeatTimer(pSyncNode);
M
Minghao Li 已提交
2299

M
Minghao Li 已提交
2300 2301
  // send heartbeat right now
  syncNodeHeartbeatPeers(pSyncNode);
M
Minghao Li 已提交
2302

2303 2304 2305 2306 2307
  // call back
  if (pSyncNode->pFsm != NULL && pSyncNode->pFsm->FpBecomeLeaderCb != NULL) {
    pSyncNode->pFsm->FpBecomeLeaderCb(pSyncNode->pFsm);
  }

M
Minghao Li 已提交
2308 2309 2310
  // min match index
  pSyncNode->minMatchIndex = SYNC_INDEX_INVALID;

M
Minghao Li 已提交
2311 2312 2313 2314 2315 2316 2317 2318 2319 2320 2321 2322 2323 2324
  // trace log
  do {
    int32_t debugStrLen = strlen(debugStr);
    if (debugStrLen < 256) {
      char eventLog[256 + 64];
      snprintf(eventLog, sizeof(eventLog), "become leader %s", debugStr);
      syncNodeEventLog(pSyncNode, eventLog);
    } else {
      char* eventLog = taosMemoryMalloc(debugStrLen + 64);
      snprintf(eventLog, debugStrLen, "become leader %s", debugStr);
      syncNodeEventLog(pSyncNode, eventLog);
      taosMemoryFree(eventLog);
    }
  } while (0);
M
Minghao Li 已提交
2325 2326 2327
}

void syncNodeCandidate2Leader(SSyncNode* pSyncNode) {
M
Minghao Li 已提交
2328 2329
  ASSERT(pSyncNode->state == TAOS_SYNC_STATE_CANDIDATE);
  ASSERT(voteGrantedMajority(pSyncNode->pVotesGranted));
2330
  syncNodeBecomeLeader(pSyncNode, "candidate to leader");
M
Minghao Li 已提交
2331

M
Minghao Li 已提交
2332 2333
  syncNodeLog2("==state change syncNodeCandidate2Leader==", pSyncNode);

M
Minghao Li 已提交
2334
  // Raft 3.6.2 Committing entries from previous terms
2335 2336
  syncNodeAppendNoop(pSyncNode);
  syncMaybeAdvanceCommitIndex(pSyncNode);
2337 2338

  if (pSyncNode->replicaNum > 1) {
M
Minghao Li 已提交
2339
    syncNodeReplicate(pSyncNode);
2340
  }
M
Minghao Li 已提交
2341 2342
}

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

M
Minghao Li 已提交
2345 2346 2347 2348 2349 2350 2351
int32_t syncNodePeerStateInit(SSyncNode* pSyncNode) {
  for (int i = 0; i < TSDB_MAX_REPLICA; ++i) {
    pSyncNode->peerStates[i].lastSendIndex = SYNC_INDEX_INVALID;
    pSyncNode->peerStates[i].lastSendTime = 0;
  }

  return 0;
M
Minghao Li 已提交
2352 2353 2354
}

void syncNodeFollower2Candidate(SSyncNode* pSyncNode) {
M
Minghao Li 已提交
2355
  ASSERT(pSyncNode->state == TAOS_SYNC_STATE_FOLLOWER);
M
Minghao Li 已提交
2356
  pSyncNode->state = TAOS_SYNC_STATE_CANDIDATE;
M
Minghao Li 已提交
2357

M
Minghao Li 已提交
2358
  syncNodeEventLog(pSyncNode, "follower to candidate");
M
Minghao Li 已提交
2359 2360 2361
}

void syncNodeLeader2Follower(SSyncNode* pSyncNode) {
M
Minghao Li 已提交
2362
  ASSERT(pSyncNode->state == TAOS_SYNC_STATE_LEADER);
2363
  syncNodeBecomeFollower(pSyncNode, "leader to follower");
M
Minghao Li 已提交
2364

M
Minghao Li 已提交
2365
  syncNodeEventLog(pSyncNode, "leader to follower");
M
Minghao Li 已提交
2366 2367 2368
}

void syncNodeCandidate2Follower(SSyncNode* pSyncNode) {
M
Minghao Li 已提交
2369
  ASSERT(pSyncNode->state == TAOS_SYNC_STATE_CANDIDATE);
2370
  syncNodeBecomeFollower(pSyncNode, "candidate to follower");
M
Minghao Li 已提交
2371

M
Minghao Li 已提交
2372
  syncNodeEventLog(pSyncNode, "candidate to follower");
M
Minghao Li 已提交
2373 2374 2375
}

// raft vote --------------
M
Minghao Li 已提交
2376 2377 2378

// just called by syncNodeVoteForSelf
// need assert
M
Minghao Li 已提交
2379
void syncNodeVoteForTerm(SSyncNode* pSyncNode, SyncTerm term, SRaftId* pRaftId) {
M
Minghao Li 已提交
2380 2381
  ASSERT(term == pSyncNode->pRaftStore->currentTerm);
  ASSERT(!raftStoreHasVoted(pSyncNode->pRaftStore));
M
Minghao Li 已提交
2382 2383 2384 2385

  raftStoreVote(pSyncNode->pRaftStore, pRaftId);
}

M
Minghao Li 已提交
2386
// simulate get vote from outside
M
Minghao Li 已提交
2387 2388 2389
void syncNodeVoteForSelf(SSyncNode* pSyncNode) {
  syncNodeVoteForTerm(pSyncNode, pSyncNode->pRaftStore->currentTerm, &(pSyncNode->myRaftId));

M
Minghao Li 已提交
2390
  SyncRequestVoteReply* pMsg = syncRequestVoteReplyBuild(pSyncNode->vgId);
M
Minghao Li 已提交
2391 2392 2393 2394 2395 2396 2397 2398 2399 2400
  pMsg->srcId = pSyncNode->myRaftId;
  pMsg->destId = pSyncNode->myRaftId;
  pMsg->term = pSyncNode->pRaftStore->currentTerm;
  pMsg->voteGranted = true;

  voteGrantedVote(pSyncNode->pVotesGranted, pMsg);
  votesRespondAdd(pSyncNode->pVotesRespond, pMsg);
  syncRequestVoteReplyDestroy(pMsg);
}

M
Minghao Li 已提交
2401
// snapshot --------------
M
Minghao Li 已提交
2402 2403

// return if has a snapshot
M
Minghao Li 已提交
2404 2405
bool syncNodeHasSnapshot(SSyncNode* pSyncNode) {
  bool      ret = false;
2406
  SSnapshot snapshot = {.data = NULL, .lastApplyIndex = -1, .lastApplyTerm = 0, .lastConfigIndex = -1};
2407 2408
  if (pSyncNode->pFsm->FpGetSnapshotInfo != NULL) {
    pSyncNode->pFsm->FpGetSnapshotInfo(pSyncNode->pFsm, &snapshot);
M
Minghao Li 已提交
2409 2410 2411 2412 2413 2414 2415
    if (snapshot.lastApplyIndex >= SYNC_INDEX_BEGIN) {
      ret = true;
    }
  }
  return ret;
}

M
Minghao Li 已提交
2416 2417
// return max(logLastIndex, snapshotLastIndex)
// if no snapshot and log, return -1
2418
SyncIndex syncNodeGetLastIndex(const SSyncNode* pSyncNode) {
M
Minghao Li 已提交
2419
  SSnapshot snapshot = {.data = NULL, .lastApplyIndex = -1, .lastApplyTerm = 0, .lastConfigIndex = -1};
2420 2421
  if (pSyncNode->pFsm->FpGetSnapshotInfo != NULL) {
    pSyncNode->pFsm->FpGetSnapshotInfo(pSyncNode->pFsm, &snapshot);
M
Minghao Li 已提交
2422 2423 2424 2425 2426 2427 2428
  }
  SyncIndex logLastIndex = pSyncNode->pLogStore->syncLogLastIndex(pSyncNode->pLogStore);

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

M
Minghao Li 已提交
2429 2430
// return the last term of snapshot and log
// if error, return SYNC_TERM_INVALID (by syncLogLastTerm)
M
Minghao Li 已提交
2431 2432
SyncTerm syncNodeGetLastTerm(SSyncNode* pSyncNode) {
  SyncTerm lastTerm = 0;
M
Minghao Li 已提交
2433 2434
  if (syncNodeHasSnapshot(pSyncNode)) {
    // has snapshot
2435
    SSnapshot snapshot = {.data = NULL, .lastApplyIndex = -1, .lastApplyTerm = 0, .lastConfigIndex = -1};
2436 2437
    if (pSyncNode->pFsm->FpGetSnapshotInfo != NULL) {
      pSyncNode->pFsm->FpGetSnapshotInfo(pSyncNode->pFsm, &snapshot);
M
Minghao Li 已提交
2438 2439
    }

M
Minghao Li 已提交
2440 2441 2442
    SyncIndex logLastIndex = pSyncNode->pLogStore->syncLogLastIndex(pSyncNode->pLogStore);
    if (logLastIndex > snapshot.lastApplyIndex) {
      lastTerm = pSyncNode->pLogStore->syncLogLastTerm(pSyncNode->pLogStore);
M
Minghao Li 已提交
2443 2444 2445 2446
    } else {
      lastTerm = snapshot.lastApplyTerm;
    }

M
Minghao Li 已提交
2447
  } else {
M
Minghao Li 已提交
2448 2449
    // no snapshot
    lastTerm = pSyncNode->pLogStore->syncLogLastTerm(pSyncNode->pLogStore);
2450
  }
M
Minghao Li 已提交
2451

M
Minghao Li 已提交
2452 2453 2454 2455 2456 2457 2458
  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);
2459 2460
  return 0;
}
M
Minghao Li 已提交
2461

M
Minghao Li 已提交
2462
// return append-entries first try index
M
Minghao Li 已提交
2463 2464 2465 2466 2467
SyncIndex syncNodeSyncStartIndex(SSyncNode* pSyncNode) {
  SyncIndex syncStartIndex = syncNodeGetLastIndex(pSyncNode) + 1;
  return syncStartIndex;
}

M
Minghao Li 已提交
2468 2469
// if index > 0, return index - 1
// else, return -1
2470 2471 2472 2473 2474 2475 2476 2477 2478
SyncIndex syncNodeGetPreIndex(SSyncNode* pSyncNode, SyncIndex index) {
  SyncIndex preIndex = index - 1;
  if (preIndex < SYNC_INDEX_INVALID) {
    preIndex = SYNC_INDEX_INVALID;
  }

  return preIndex;
}

M
Minghao Li 已提交
2479 2480 2481 2482
// if index < 0, return SYNC_TERM_INVALID
// if index == 0, return 0
// if index > 0, return preTerm
// if error, return SYNC_TERM_INVALID
2483 2484 2485 2486 2487 2488 2489 2490 2491 2492 2493 2494 2495
SyncTerm syncNodeGetPreTerm(SSyncNode* pSyncNode, SyncIndex index) {
  if (index < SYNC_INDEX_BEGIN) {
    return SYNC_TERM_INVALID;
  }

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

  SyncTerm        preTerm = 0;
  SyncIndex       preIndex = index - 1;
  SSyncRaftEntry* pPreEntry = NULL;
  int32_t         code = pSyncNode->pLogStore->syncLogGetEntry(pSyncNode->pLogStore, preIndex, &pPreEntry);
M
Minghao Li 已提交
2496 2497 2498 2499 2500 2501

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

2502 2503 2504 2505 2506 2507
  if (code == 0) {
    ASSERT(pPreEntry != NULL);
    preTerm = pPreEntry->term;
    taosMemoryFree(pPreEntry);
    return preTerm;
  } else {
2508 2509 2510 2511
    if (pSyncNode->pFsm->FpGetSnapshotInfo != NULL) {
      pSyncNode->pFsm->FpGetSnapshotInfo(pSyncNode->pFsm, &snapshot);
      if (snapshot.lastApplyIndex == preIndex) {
        return snapshot.lastApplyTerm;
2512 2513 2514 2515
      }
    }
  }

2516 2517
  do {
    char logBuf[128];
S
Shengliang Guan 已提交
2518 2519
    snprintf(logBuf, sizeof(logBuf),
             "sync node get pre term error, index:%" PRId64 ", snap-index:%" PRId64 ", snap-term:%" PRIu64, index,
M
Minghao Li 已提交
2520
             snapshot.lastApplyIndex, snapshot.lastApplyTerm);
2521 2522 2523
    syncNodeErrorLog(pSyncNode, logBuf);
  } while (0);

2524 2525
  return SYNC_TERM_INVALID;
}
M
Minghao Li 已提交
2526 2527 2528 2529

// 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 已提交
2530
  *pPreTerm = syncNodeGetPreTerm(pSyncNode, index);
M
Minghao Li 已提交
2531 2532 2533
  return 0;
}

M
Minghao Li 已提交
2534 2535 2536
// for debug --------------
void syncNodePrint(SSyncNode* pObj) {
  char* serialized = syncNode2Str(pObj);
S
Shengliang Guan 已提交
2537
  printf("syncNodePrint | len:%d | %s \n", (int32_t)strlen(serialized), serialized);
M
Minghao Li 已提交
2538
  fflush(NULL);
wafwerar's avatar
wafwerar 已提交
2539
  taosMemoryFree(serialized);
M
Minghao Li 已提交
2540 2541 2542 2543
}

void syncNodePrint2(char* s, SSyncNode* pObj) {
  char* serialized = syncNode2Str(pObj);
S
Shengliang Guan 已提交
2544
  printf("syncNodePrint2 | len:%d | %s | %s \n", (int32_t)strlen(serialized), s, serialized);
M
Minghao Li 已提交
2545
  fflush(NULL);
wafwerar's avatar
wafwerar 已提交
2546
  taosMemoryFree(serialized);
M
Minghao Li 已提交
2547 2548 2549 2550
}

void syncNodeLog(SSyncNode* pObj) {
  char* serialized = syncNode2Str(pObj);
S
Shengliang Guan 已提交
2551
  sTraceLong("syncNodeLog | len:%d | %s", (int32_t)strlen(serialized), serialized);
wafwerar's avatar
wafwerar 已提交
2552
  taosMemoryFree(serialized);
M
Minghao Li 已提交
2553 2554 2555
}

void syncNodeLog2(char* s, SSyncNode* pObj) {
2556 2557
  if (gRaftDetailLog) {
    char* serialized = syncNode2Str(pObj);
S
Shengliang Guan 已提交
2558
    sTraceLong("syncNodeLog2 | len:%d | %s | %s", (int32_t)strlen(serialized), s, serialized);
2559 2560
    taosMemoryFree(serialized);
  }
M
Minghao Li 已提交
2561 2562
}

M
Minghao Li 已提交
2563 2564
void syncNodeLog3(char* s, SSyncNode* pObj) {
  char* serialized = syncNode2Str(pObj);
S
Shengliang Guan 已提交
2565
  sTraceLong("syncNodeLog3 | len:%d | %s | %s", (int32_t)strlen(serialized), s, serialized);
M
Minghao Li 已提交
2566 2567 2568
  taosMemoryFree(serialized);
}

M
Minghao Li 已提交
2569
// ------ local funciton ---------
M
Minghao Li 已提交
2570
// enqueue message ----
M
Minghao Li 已提交
2571 2572
static void syncNodeEqPingTimer(void* param, void* tmrId) {
  SSyncNode* pSyncNode = (SSyncNode*)param;
M
Minghao Li 已提交
2573
  if (atomic_load_64(&pSyncNode->pingTimerLogicClockUser) <= atomic_load_64(&pSyncNode->pingTimerLogicClock)) {
M
Minghao Li 已提交
2574
    SyncTimeout* pSyncMsg = syncTimeoutBuild2(SYNC_TIMEOUT_PING, atomic_load_64(&pSyncNode->pingTimerLogicClock),
M
Minghao Li 已提交
2575
                                              pSyncNode->pingTimerMS, pSyncNode->vgId, pSyncNode);
M
Minghao Li 已提交
2576 2577
    SRpcMsg      rpcMsg;
    syncTimeout2RpcMsg(pSyncMsg, &rpcMsg);
M
Minghao Li 已提交
2578
    syncRpcMsgLog2((char*)"==syncNodeEqPingTimer==", &rpcMsg);
S
Shengliang Guan 已提交
2579 2580
    if (pSyncNode->syncEqMsg != NULL) {
      int32_t code = pSyncNode->syncEqMsg(pSyncNode->msgcb, &rpcMsg);
2581
      if (code != 0) {
S
Shengliang Guan 已提交
2582
        sError("vgId:%d, sync enqueue ping msg error, code:%d", pSyncNode->vgId, code);
2583 2584 2585 2586
        rpcFreeCont(rpcMsg.pCont);
        syncTimeoutDestroy(pSyncMsg);
        return;
      }
M
Minghao Li 已提交
2587
    } else {
S
Shengliang Guan 已提交
2588
      sTrace("syncNodeEqPingTimer pSyncNode->syncEqMsg is NULL");
M
Minghao Li 已提交
2589
    }
M
Minghao Li 已提交
2590 2591
    syncTimeoutDestroy(pSyncMsg);

S
Shengliang Guan 已提交
2592 2593
    if (syncIsInit()) {
      taosTmrReset(syncNodeEqPingTimer, pSyncNode->pingTimerMS, pSyncNode, syncEnv()->pTimerManager,
2594 2595 2596 2597 2598
                   &pSyncNode->pPingTimer);
    } else {
      sError("sync env is stop, syncNodeEqPingTimer");
    }

M
Minghao Li 已提交
2599
  } else {
S
Shengliang Guan 已提交
2600
    sTrace("==syncNodeEqPingTimer== pingTimerLogicClock:%" PRIu64 ", pingTimerLogicClockUser:%" PRIu64,
M
Minghao Li 已提交
2601
           pSyncNode->pingTimerLogicClock, pSyncNode->pingTimerLogicClockUser);
M
Minghao Li 已提交
2602 2603 2604 2605
  }
}

static void syncNodeEqElectTimer(void* param, void* tmrId) {
M
Minghao Li 已提交
2606 2607
  SElectTimer* pElectTimer = (SElectTimer*)param;
  SSyncNode*   pSyncNode = pElectTimer->pSyncNode;
M
Minghao Li 已提交
2608

M
Minghao Li 已提交
2609 2610
  SyncTimeout* pSyncMsg = syncTimeoutBuild2(SYNC_TIMEOUT_ELECTION, pElectTimer->logicClock, pSyncNode->electTimerMS,
                                            pSyncNode->vgId, pSyncNode);
M
Minghao Li 已提交
2611 2612
  SRpcMsg      rpcMsg;
  syncTimeout2RpcMsg(pSyncMsg, &rpcMsg);
2613
  if (pSyncNode->syncEqMsg != NULL && pSyncNode->msgcb != NULL && pSyncNode->msgcb->putToQueueFp != NULL) {
S
Shengliang Guan 已提交
2614
    int32_t code = pSyncNode->syncEqMsg(pSyncNode->msgcb, &rpcMsg);
M
Minghao Li 已提交
2615 2616 2617 2618
    if (code != 0) {
      sError("vgId:%d, sync enqueue elect msg error, code:%d", pSyncNode->vgId, code);
      rpcFreeCont(rpcMsg.pCont);
      syncTimeoutDestroy(pSyncMsg);
M
Minghao Li 已提交
2619
      taosMemoryFree(pElectTimer);
M
Minghao Li 已提交
2620
      return;
2621
    }
M
Minghao Li 已提交
2622 2623 2624

    do {
      char logBuf[128];
M
Minghao Li 已提交
2625
      snprintf(logBuf, sizeof(logBuf), "eq elect timer lc:%" PRIu64, pSyncMsg->logicClock);
M
Minghao Li 已提交
2626
      syncNodeEventLog(pSyncNode, logBuf);
M
Minghao Li 已提交
2627 2628
    } while (0);

M
Minghao Li 已提交
2629
  } else {
S
Shengliang Guan 已提交
2630
    sTrace("syncNodeEqElectTimer syncEqMsg is NULL");
M
Minghao Li 已提交
2631
  }
M
Minghao Li 已提交
2632

M
Minghao Li 已提交
2633
  syncTimeoutDestroy(pSyncMsg);
M
Minghao Li 已提交
2634
  taosMemoryFree(pElectTimer);
M
Minghao Li 已提交
2635

M
Minghao Li 已提交
2636
#if 0
M
Minghao Li 已提交
2637
  // reset timer ms
S
Shengliang Guan 已提交
2638
  if (syncIsInit() && pSyncNode->electBaseLine > 0) {
M
Minghao Li 已提交
2639
    pSyncNode->electTimerMS = syncUtilElectRandomMS(pSyncNode->electBaseLine, 2 * pSyncNode->electBaseLine);
S
Shengliang Guan 已提交
2640
    taosTmrReset(syncNodeEqElectTimer, pSyncNode->electTimerMS, pSyncNode, syncEnv()->pTimerManager,
M
Minghao Li 已提交
2641 2642 2643
                 &pSyncNode->pElectTimer);
  } else {
    sError("sync env is stop, syncNodeEqElectTimer");
M
Minghao Li 已提交
2644
  }
M
Minghao Li 已提交
2645
#endif
M
Minghao Li 已提交
2646 2647
}

M
Minghao Li 已提交
2648 2649
static void syncNodeEqHeartbeatTimer(void* param, void* tmrId) {
  SSyncNode* pSyncNode = (SSyncNode*)param;
2650 2651 2652

  syncNodeEventLog(pSyncNode, "eq hb timer");

2653 2654 2655 2656 2657 2658 2659 2660 2661
  if (pSyncNode->replicaNum > 1) {
    if (atomic_load_64(&pSyncNode->heartbeatTimerLogicClockUser) <=
        atomic_load_64(&pSyncNode->heartbeatTimerLogicClock)) {
      SyncTimeout* pSyncMsg =
          syncTimeoutBuild2(SYNC_TIMEOUT_HEARTBEAT, atomic_load_64(&pSyncNode->heartbeatTimerLogicClock),
                            pSyncNode->heartbeatTimerMS, pSyncNode->vgId, pSyncNode);
      SRpcMsg rpcMsg;
      syncTimeout2RpcMsg(pSyncMsg, &rpcMsg);
      syncRpcMsgLog2((char*)"==syncNodeEqHeartbeatTimer==", &rpcMsg);
S
Shengliang Guan 已提交
2662 2663
      if (pSyncNode->syncEqMsg != NULL) {
        int32_t code = pSyncNode->syncEqMsg(pSyncNode->msgcb, &rpcMsg);
2664
        if (code != 0) {
S
Shengliang Guan 已提交
2665
          sError("vgId:%d, sync enqueue timer msg error, code:%d", pSyncNode->vgId, code);
2666 2667 2668 2669 2670
          rpcFreeCont(rpcMsg.pCont);
          syncTimeoutDestroy(pSyncMsg);
          return;
        }
      } else {
S
Shengliang Guan 已提交
2671
        sError("vgId:%d, enqueue msg cb ptr (i.e. syncEqMsg) not set.", pSyncNode->vgId);
2672
      }
2673
      syncTimeoutDestroy(pSyncMsg);
M
Minghao Li 已提交
2674

S
Shengliang Guan 已提交
2675 2676
      if (syncIsInit()) {
        taosTmrReset(syncNodeEqHeartbeatTimer, pSyncNode->heartbeatTimerMS, pSyncNode, syncEnv()->pTimerManager,
2677 2678 2679 2680
                     &pSyncNode->pHeartbeatTimer);
      } else {
        sError("sync env is stop, syncNodeEqHeartbeatTimer");
      }
2681
    } else {
2682 2683 2684
      sTrace("==syncNodeEqHeartbeatTimer== heartbeatTimerLogicClock:%" PRIu64 ", heartbeatTimerLogicClockUser:%" PRIu64
             "",
             pSyncNode->heartbeatTimerLogicClock, pSyncNode->heartbeatTimerLogicClockUser);
2685
    }
M
Minghao Li 已提交
2686 2687 2688
  }
}

2689 2690 2691 2692 2693
static void syncNodeEqPeerHeartbeatTimer(void* param, void* tmrId) {
  SSyncHbTimerData* pData = (SSyncHbTimerData*)param;
  SSyncNode*        pSyncNode = pData->pSyncNode;
  SSyncTimer*       pSyncTimer = pData->pTimer;

M
Minghao Li 已提交
2694 2695 2696 2697
  if (pSyncNode->state != TAOS_SYNC_STATE_LEADER) {
    return;
  }

S
Shengliang Guan 已提交
2698
  // syncNodeEventLog(pSyncNode, "eq peer hb timer");
2699 2700 2701 2702 2703 2704 2705 2706 2707 2708 2709

  int64_t timerLogicClock = atomic_load_64(&pSyncTimer->logicClock);
  int64_t msgLogicClock = atomic_load_64(&pData->logicClock);

  if (pSyncNode->replicaNum > 1) {
    if (timerLogicClock == msgLogicClock) {
      SyncHeartbeat* pSyncMsg = syncHeartbeatBuild(pSyncNode->vgId);
      pSyncMsg->srcId = pSyncNode->myRaftId;
      pSyncMsg->destId = pData->destId;
      pSyncMsg->term = pSyncNode->pRaftStore->currentTerm;
      pSyncMsg->commitIndex = pSyncNode->commitIndex;
M
Minghao Li 已提交
2710
      pSyncMsg->minMatchIndex = syncMinMatchIndex(pSyncNode);
2711 2712 2713 2714 2715 2716 2717
      pSyncMsg->privateTerm = 0;

      SRpcMsg rpcMsg;
      syncHeartbeat2RpcMsg(pSyncMsg, &rpcMsg);

// eq msg
#if 0
S
Shengliang Guan 已提交
2718 2719
      if (pSyncNode->syncEqCtrlMsg != NULL) {
        int32_t code = pSyncNode->syncEqCtrlMsg(pSyncNode->msgcb, &rpcMsg);
2720 2721 2722 2723 2724 2725 2726
        if (code != 0) {
          sError("vgId:%d, sync ctrl enqueue timer msg error, code:%d", pSyncNode->vgId, code);
          rpcFreeCont(rpcMsg.pCont);
          syncHeartbeatDestroy(pSyncMsg);
          return;
        }
      } else {
S
Shengliang Guan 已提交
2727
        sError("vgId:%d, enqueue ctrl msg cb ptr (i.e. syncEqMsg) not set.", pSyncNode->vgId);
2728 2729 2730 2731
      }
#endif

      // send msg
M
Minghao Li 已提交
2732
      syncNodeSendHeartbeat(pSyncNode, &(pSyncMsg->destId), pSyncMsg);
2733 2734 2735

      syncHeartbeatDestroy(pSyncMsg);

S
Shengliang Guan 已提交
2736 2737
      if (syncIsInit()) {
        taosTmrReset(syncNodeEqPeerHeartbeatTimer, pSyncTimer->timerMS, pData, syncEnv()->pTimerManager,
2738 2739 2740 2741 2742 2743 2744 2745 2746 2747 2748 2749
                     &pSyncTimer->pTimer);
      } else {
        sError("sync env is stop, syncNodeEqHeartbeatTimer");
      }

    } else {
      sTrace("==syncNodeEqPeerHeartbeatTimer== timerLogicClock:%" PRIu64 ", msgLogicClock:%" PRIu64 "", timerLogicClock,
             msgLogicClock);
    }
  }
}

M
Minghao Li 已提交
2750 2751
static int32_t syncNodeEqNoop(SSyncNode* ths) {
  int32_t ret = 0;
M
Minghao Li 已提交
2752
  ASSERT(ths->state == TAOS_SYNC_STATE_LEADER);
M
Minghao Li 已提交
2753

2754
  SyncIndex       index = ths->pLogStore->syncLogWriteIndex(ths->pLogStore);
M
Minghao Li 已提交
2755
  SyncTerm        term = ths->pRaftStore->currentTerm;
M
Minghao Li 已提交
2756
  SSyncRaftEntry* pEntry = syncEntryBuildNoop(term, index, ths->vgId);
M
Minghao Li 已提交
2757
  ASSERT(pEntry != NULL);
M
Minghao Li 已提交
2758 2759 2760 2761

  uint32_t           entryLen;
  char*              serialized = syncEntrySerialize(pEntry, &entryLen);
  SyncClientRequest* pSyncMsg = syncClientRequestBuild(entryLen);
M
Minghao Li 已提交
2762
  ASSERT(pSyncMsg->dataLen == entryLen);
M
Minghao Li 已提交
2763 2764
  memcpy(pSyncMsg->data, serialized, entryLen);

S
Shengliang Guan 已提交
2765
  SRpcMsg rpcMsg = {0};
M
Minghao Li 已提交
2766
  syncClientRequest2RpcMsg(pSyncMsg, &rpcMsg);
S
Shengliang Guan 已提交
2767 2768
  if (ths->syncEqMsg != NULL) {
    ths->syncEqMsg(ths->msgcb, &rpcMsg);
M
Minghao Li 已提交
2769
  } else {
S
Shengliang Guan 已提交
2770
    sTrace("syncNodeEqNoop pSyncNode->syncEqMsg is NULL");
M
Minghao Li 已提交
2771
  }
M
Minghao Li 已提交
2772

M
Minghao Li 已提交
2773
  syncEntryDestory(pEntry);
wafwerar's avatar
wafwerar 已提交
2774
  taosMemoryFree(serialized);
M
Minghao Li 已提交
2775 2776 2777 2778 2779
  syncClientRequestDestroy(pSyncMsg);

  return ret;
}

2780 2781 2782 2783 2784 2785 2786 2787 2788 2789 2790 2791 2792 2793
static void deleteCacheEntry(const void* key, size_t keyLen, void* value) { taosMemoryFree(value); }

static int32_t syncCacheEntry(SSyncLogStore* pLogStore, SSyncRaftEntry* pEntry, LRUHandle** h) {
  int       code = 0;
  int       entryLen = sizeof(*pEntry) + pEntry->dataLen;
  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;
}

M
Minghao Li 已提交
2794 2795 2796
static int32_t syncNodeAppendNoop(SSyncNode* ths) {
  int32_t ret = 0;

2797
  SyncIndex       index = ths->pLogStore->syncLogWriteIndex(ths->pLogStore);
M
Minghao Li 已提交
2798
  SyncTerm        term = ths->pRaftStore->currentTerm;
M
Minghao Li 已提交
2799
  SSyncRaftEntry* pEntry = syncEntryBuildNoop(term, index, ths->vgId);
M
Minghao Li 已提交
2800
  ASSERT(pEntry != NULL);
M
Minghao Li 已提交
2801

2802 2803 2804
  LRUHandle* h = NULL;
  syncCacheEntry(ths->pLogStore, pEntry, &h);

M
Minghao Li 已提交
2805
  if (ths->state == TAOS_SYNC_STATE_LEADER) {
2806
    int32_t code = ths->pLogStore->syncLogAppendEntry(ths->pLogStore, pEntry);
2807
    if (code != 0) {
M
Minghao Li 已提交
2808
      syncNodeErrorLog(ths, "append noop error");
2809 2810
      return -1;
    }
M
Minghao Li 已提交
2811 2812
  }

2813 2814 2815 2816 2817 2818
  if (h) {
    taosLRUCacheRelease(ths->pLogStore->pCache, h, false);
  } else {
    syncEntryDestory(pEntry);
  }

M
Minghao Li 已提交
2819 2820 2821
  return ret;
}

M
Minghao Li 已提交
2822
// on message ----
M
Minghao Li 已提交
2823 2824 2825
int32_t syncNodeOnPing(SSyncNode* ths, SyncPing* pMsg) {
  sTrace("vgId:%d, recv sync-ping", ths->vgId);

M
Minghao Li 已提交
2826
  SyncPingReply* pMsgReply = syncPingReplyBuild3(&ths->myRaftId, &pMsg->srcId, ths->vgId);
M
Minghao Li 已提交
2827 2828
  SRpcMsg        rpcMsg;
  syncPingReply2RpcMsg(pMsgReply, &rpcMsg);
M
Minghao Li 已提交
2829 2830

  /*
M
Minghao Li 已提交
2831 2832 2833 2834 2835
    // htonl
    SMsgHead* pHead = rpcMsg.pCont;
    pHead->contLen = htonl(pHead->contLen);
    pHead->vgId = htonl(pHead->vgId);
  */
M
Minghao Li 已提交
2836

M
Minghao Li 已提交
2837
  syncNodeSendMsgById(&pMsgReply->destId, ths, &rpcMsg);
M
Minghao Li 已提交
2838
  syncPingReplyDestroy(pMsgReply);
M
Minghao Li 已提交
2839

M
Minghao Li 已提交
2840
  return 0;
M
Minghao Li 已提交
2841 2842
}

M
Minghao Li 已提交
2843
int32_t syncNodeOnPingReply(SSyncNode* ths, SyncPingReply* pMsg) {
M
Minghao Li 已提交
2844
  int32_t ret = 0;
M
Minghao Li 已提交
2845
  sTrace("vgId:%d, recv sync-ping-reply", ths->vgId);
M
Minghao Li 已提交
2846 2847
  return ret;
}
M
Minghao Li 已提交
2848

2849 2850 2851 2852 2853 2854 2855 2856 2857 2858 2859 2860
int32_t syncNodeOnHeartbeat(SSyncNode* ths, SyncHeartbeat* pMsg) {
  syncLogRecvHeartbeat(ths, pMsg, "");

  SyncHeartbeatReply* pMsgReply = syncHeartbeatReplyBuild(ths->vgId);
  pMsgReply->destId = pMsg->srcId;
  pMsgReply->srcId = ths->myRaftId;
  pMsgReply->term = ths->pRaftStore->currentTerm;
  pMsgReply->privateTerm = 8864;  // magic number

  SRpcMsg rpcMsg;
  syncHeartbeatReply2RpcMsg(pMsgReply, &rpcMsg);

M
Minghao Li 已提交
2861
  if (pMsg->term == ths->pRaftStore->currentTerm && ths->state != TAOS_SYNC_STATE_LEADER) {
2862
    syncNodeResetElectTimer(ths);
M
Minghao Li 已提交
2863
    ths->minMatchIndex = pMsg->minMatchIndex;
2864 2865 2866 2867 2868 2869 2870 2871

#if 0
    if (ths->state == TAOS_SYNC_STATE_FOLLOWER) {
      syncNodeFollowerCommit(ths, pMsg->commitIndex);
    }
#endif
  }

M
Minghao Li 已提交
2872
  if (pMsg->term >= ths->pRaftStore->currentTerm && ths->state != TAOS_SYNC_STATE_FOLLOWER) {
2873 2874 2875 2876 2877 2878 2879 2880
    // syncNodeStepDown(ths, pMsg->term);
    SyncLocalCmd* pSyncMsg = syncLocalCmdBuild(ths->vgId);
    pSyncMsg->cmd = SYNC_LOCAL_CMD_STEP_DOWN;
    pSyncMsg->sdNewTerm = pMsg->term;

    SRpcMsg rpcMsgLocalCmd;
    syncLocalCmd2RpcMsg(pSyncMsg, &rpcMsgLocalCmd);

S
Shengliang Guan 已提交
2881 2882
    if (ths->syncEqMsg != NULL && ths->msgcb != NULL) {
      int32_t code = ths->syncEqMsg(ths->msgcb, &rpcMsgLocalCmd);
2883 2884 2885 2886 2887 2888 2889 2890 2891
      if (code != 0) {
        sError("vgId:%d, sync enqueue step-down msg error, code:%d", ths->vgId, code);
        rpcFreeCont(rpcMsgLocalCmd.pCont);
      } else {
        sTrace("vgId:%d, sync enqueue step-down msg, new-term: %" PRIu64, ths->vgId, pSyncMsg->sdNewTerm);
      }
    }

    syncLocalCmdDestroy(pSyncMsg);
M
Minghao Li 已提交
2892 2893
  }

2894 2895 2896 2897 2898 2899 2900 2901 2902
  /*
    // htonl
    SMsgHead* pHead = rpcMsg.pCont;
    pHead->contLen = htonl(pHead->contLen);
    pHead->vgId = htonl(pHead->vgId);
  */

  // reply
  syncNodeSendMsgById(&pMsgReply->destId, ths, &rpcMsg);
M
Minghao Li 已提交
2903
  syncHeartbeatReplyDestroy(pMsgReply);
2904 2905 2906 2907 2908 2909 2910 2911 2912 2913 2914 2915 2916

  return 0;
}

int32_t syncNodeOnHeartbeatReply(SSyncNode* ths, SyncHeartbeatReply* pMsg) {
  syncLogRecvHeartbeatReply(ths, pMsg, "");

  // update last reply time, make decision whether the other node is alive or not
  syncIndexMgrSetRecvTime(ths->pMatchIndex, &(pMsg->destId), pMsg->startTime);

  return 0;
}

M
Minghao Li 已提交
2917
int32_t syncNodeOnLocalCmd(SSyncNode* ths, SyncLocalCmd* pMsg) {
2918 2919
  syncLogRecvLocalCmd(ths, pMsg, "");

M
Minghao Li 已提交
2920 2921 2922 2923 2924 2925 2926 2927 2928 2929
  if (pMsg->cmd == SYNC_LOCAL_CMD_STEP_DOWN) {
    syncNodeStepDown(ths, pMsg->sdNewTerm);

  } else {
    syncNodeErrorLog(ths, "error local cmd");
  }

  return 0;
}

M
Minghao Li 已提交
2930 2931 2932 2933 2934 2935 2936 2937 2938 2939
// 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 已提交
2940

M
Minghao Li 已提交
2941
int32_t syncNodeOnClientRequest(SSyncNode* ths, SyncClientRequest* pMsg, SyncIndex* pRetIndex) {
2942 2943
  syncNodeEventLog(ths, "on client request");

M
Minghao Li 已提交
2944
  int32_t ret = 0;
2945
  int32_t code = 0;
M
Minghao Li 已提交
2946

M
Minghao Li 已提交
2947
  SyncIndex       index = ths->pLogStore->syncLogWriteIndex(ths->pLogStore);
M
Minghao Li 已提交
2948 2949
  SyncTerm        term = ths->pRaftStore->currentTerm;
  SSyncRaftEntry* pEntry = syncEntryBuild2((SyncClientRequest*)pMsg, term, index);
M
Minghao Li 已提交
2950
  ASSERT(pEntry != NULL);
M
Minghao Li 已提交
2951

2952 2953 2954
  LRUHandle* h = NULL;
  syncCacheEntry(ths->pLogStore, pEntry, &h);

M
Minghao Li 已提交
2955
  if (ths->state == TAOS_SYNC_STATE_LEADER) {
2956 2957 2958 2959
    // append entry
    code = ths->pLogStore->syncLogAppendEntry(ths->pLogStore, pEntry);
    if (code != 0) {
      // del resp mgr, call FpCommitCb
M
Minghao Li 已提交
2960
      ASSERT(0);
2961 2962
      return -1;
    }
M
Minghao Li 已提交
2963

2964 2965
    // if mulit replica, start replicate right now
    if (ths->replicaNum > 1) {
M
Minghao Li 已提交
2966
      syncNodeReplicate(ths);
2967
    }
2968

2969 2970 2971 2972
    // if only myself, maybe commit right now
    if (ths->replicaNum == 1) {
      syncMaybeAdvanceCommitIndex(ths);
    }
M
Minghao Li 已提交
2973 2974
  }

2975 2976 2977 2978 2979 2980 2981 2982
  if (pRetIndex != NULL) {
    if (ret == 0 && pEntry != NULL) {
      *pRetIndex = pEntry->index;
    } else {
      *pRetIndex = SYNC_INDEX_INVALID;
    }
  }

2983 2984 2985 2986 2987 2988
  if (h) {
    taosLRUCacheRelease(ths->pLogStore->pCache, h, false);
  } else {
    syncEntryDestory(pEntry);
  }

M
Minghao Li 已提交
2989
  return ret;
2990
}
M
Minghao Li 已提交
2991

S
Shengliang Guan 已提交
2992 2993 2994
const char* syncStr(ESyncState state) {
  switch (state) {
    case TAOS_SYNC_STATE_FOLLOWER:
2995
      return "follower";
S
Shengliang Guan 已提交
2996
    case TAOS_SYNC_STATE_CANDIDATE:
2997
      return "candidate";
S
Shengliang Guan 已提交
2998
    case TAOS_SYNC_STATE_LEADER:
2999
      return "leader";
S
Shengliang Guan 已提交
3000
    default:
3001
      return "error";
S
Shengliang Guan 已提交
3002
  }
M
Minghao Li 已提交
3003
}
3004

3005
int32_t syncDoLeaderTransfer(SSyncNode* ths, SRpcMsg* pRpcMsg, SSyncRaftEntry* pEntry) {
3006 3007 3008 3009
  if (ths->state != TAOS_SYNC_STATE_FOLLOWER) {
    syncNodeEventLog(ths, "I am not follower, can not do leader transfer");
    return 0;
  }
3010 3011 3012 3013 3014 3015

  if (!ths->restoreFinish) {
    syncNodeEventLog(ths, "restore not finish, can not do leader transfer");
    return 0;
  }

3016 3017
  if (pEntry->term < ths->pRaftStore->currentTerm) {
    char logBuf[128];
S
Shengliang Guan 已提交
3018
    snprintf(logBuf, sizeof(logBuf), "little term:%" PRIu64 ", can not do leader transfer", pEntry->term);
3019 3020 3021 3022 3023 3024
    syncNodeEventLog(ths, logBuf);
    return 0;
  }

  if (pEntry->index < syncNodeGetLastIndex(ths)) {
    char logBuf[128];
S
Shengliang Guan 已提交
3025
    snprintf(logBuf, sizeof(logBuf), "little index:%" PRId64 ", can not do leader transfer", pEntry->index);
3026
    syncNodeEventLog(ths, logBuf);
3027 3028 3029
    return 0;
  }

3030 3031 3032 3033 3034 3035 3036
  /*
    if (ths->vgId > 1) {
      syncNodeEventLog(ths, "I am vnode, can not do leader transfer");
      return 0;
    }
  */

M
Minghao Li 已提交
3037 3038
  SyncLeaderTransfer* pSyncLeaderTransfer = syncLeaderTransferFromRpcMsg2(pRpcMsg);

3039 3040
  do {
    char logBuf[128];
S
Shengliang Guan 已提交
3041
    snprintf(logBuf, sizeof(logBuf), "do leader transfer, index:%" PRId64, pEntry->index);
3042 3043
    syncNodeEventLog(ths, logBuf);
  } while (0);
M
Minghao Li 已提交
3044

M
Minghao Li 已提交
3045 3046 3047
  bool sameId = syncUtilSameId(&(pSyncLeaderTransfer->newLeaderId), &(ths->myRaftId));
  bool sameNodeInfo = strcmp(pSyncLeaderTransfer->newNodeInfo.nodeFqdn, ths->myNodeInfo.nodeFqdn) == 0 &&
                      pSyncLeaderTransfer->newNodeInfo.nodePort == ths->myNodeInfo.nodePort;
M
Minghao Li 已提交
3048

M
Minghao Li 已提交
3049 3050
  bool same = sameId || sameNodeInfo;
  if (same) {
M
Minghao Li 已提交
3051 3052 3053 3054
    // reset elect timer now!
    int32_t electMS = 1;
    int32_t ret = syncNodeRestartElectTimer(ths, electMS);
    ASSERT(ret == 0);
M
Minghao Li 已提交
3055 3056

    char eventLog[256];
S
Shengliang Guan 已提交
3057
    snprintf(eventLog, sizeof(eventLog), "maybe leader transfer to %s:%d %" PRIu64,
M
Minghao Li 已提交
3058 3059 3060
             pSyncLeaderTransfer->newNodeInfo.nodeFqdn, pSyncLeaderTransfer->newNodeInfo.nodePort,
             pSyncLeaderTransfer->newLeaderId.addr);
    syncNodeEventLog(ths, eventLog);
3061 3062
  }

M
Minghao Li 已提交
3063
  if (ths->pFsm->FpLeaderTransferCb != NULL) {
S
Shengliang Guan 已提交
3064 3065 3066 3067 3068 3069 3070 3071 3072 3073 3074 3075
    SFsmCbMeta cbMeta = {
        cbMeta.code = 0,
        cbMeta.currentTerm = ths->pRaftStore->currentTerm,
        cbMeta.flag = 0,
        cbMeta.index = pEntry->index,
        cbMeta.lastConfigIndex = syncNodeGetSnapshotConfigIndex(ths, pEntry->index),
        cbMeta.isWeak = pEntry->isWeak,
        cbMeta.seqNum = pEntry->seqNum,
        cbMeta.state = ths->state,
        cbMeta.term = pEntry->term,
    };
    ths->pFsm->FpLeaderTransferCb(ths->pFsm, pRpcMsg, &cbMeta);
3076 3077
  }

M
Minghao Li 已提交
3078
  syncLeaderTransferDestroy(pSyncLeaderTransfer);
3079 3080 3081
  return 0;
}

3082 3083 3084 3085 3086 3087 3088 3089 3090 3091 3092 3093 3094 3095 3096
int32_t syncNodeUpdateNewConfigIndex(SSyncNode* ths, SSyncCfg* pNewCfg) {
  for (int i = 0; i < pNewCfg->replicaNum; ++i) {
    SRaftId raftId;
    raftId.addr = syncUtilAddr2U64((pNewCfg->nodeInfo)[i].nodeFqdn, (pNewCfg->nodeInfo)[i].nodePort);
    raftId.vgId = ths->vgId;

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

  return -1;
}

M
Minghao Li 已提交
3097 3098 3099 3100 3101 3102 3103 3104 3105 3106 3107 3108 3109 3110 3111 3112 3113 3114 3115 3116 3117 3118
static int32_t syncNodeConfigChangeFinish(SSyncNode* ths, SRpcMsg* pRpcMsg, SSyncRaftEntry* pEntry) {
  SyncReconfigFinish* pFinish = syncReconfigFinishFromRpcMsg2(pRpcMsg);
  ASSERT(pFinish);

  if (ths->pFsm->FpReConfigCb != NULL) {
    SReConfigCbMeta cbMeta = {0};
    cbMeta.code = 0;
    cbMeta.index = pEntry->index;
    cbMeta.term = pEntry->term;
    cbMeta.seqNum = pEntry->seqNum;
    cbMeta.lastConfigIndex = syncNodeGetSnapshotConfigIndex(ths, pEntry->index);
    cbMeta.state = ths->state;
    cbMeta.currentTerm = ths->pRaftStore->currentTerm;
    cbMeta.isWeak = pEntry->isWeak;
    cbMeta.flag = 0;

    cbMeta.oldCfg = pFinish->oldCfg;
    cbMeta.newCfg = pFinish->newCfg;
    cbMeta.newCfgIndex = pFinish->newCfgIndex;
    cbMeta.newCfgTerm = pFinish->newCfgTerm;
    cbMeta.newCfgSeqNum = pFinish->newCfgSeqNum;

S
Shengliang Guan 已提交
3119
    ths->pFsm->FpReConfigCb(ths->pFsm, pRpcMsg, &cbMeta);
M
Minghao Li 已提交
3120 3121
  }

3122
  // clear changing
M
Minghao Li 已提交
3123 3124 3125 3126 3127
  ths->changing = false;

  char  tmpbuf[512];
  char* oldStr = syncCfg2SimpleStr(&(pFinish->oldCfg));
  char* newStr = syncCfg2SimpleStr(&(pFinish->newCfg));
S
Shengliang Guan 已提交
3128
  snprintf(tmpbuf, sizeof(tmpbuf), "config change finish from %d to %d, index:%" PRId64 ", %s  -->  %s",
M
Minghao Li 已提交
3129 3130 3131 3132 3133 3134 3135 3136 3137 3138 3139 3140
           pFinish->oldCfg.replicaNum, pFinish->newCfg.replicaNum, pFinish->newCfgIndex, oldStr, newStr);
  taosMemoryFree(oldStr);
  taosMemoryFree(newStr);
  syncNodeEventLog(ths, tmpbuf);

  syncReconfigFinishDestroy(pFinish);

  return 0;
}

static int32_t syncNodeConfigChange(SSyncNode* ths, SRpcMsg* pRpcMsg, SSyncRaftEntry* pEntry,
                                    SyncReconfigFinish* pFinish) {
3141 3142 3143
  // set changing
  ths->changing = true;

M
Minghao Li 已提交
3144
  // old config
3145 3146
  SSyncCfg oldSyncCfg = ths->pRaftCfg->cfg;

M
Minghao Li 已提交
3147
  // new config
3148 3149 3150 3151 3152
  SSyncCfg newSyncCfg;
  int32_t  ret = syncCfgFromStr(pRpcMsg->pCont, &newSyncCfg);
  ASSERT(ret == 0);

  // update new config myIndex
3153 3154
  syncNodeUpdateNewConfigIndex(ths, &newSyncCfg);

M
Minghao Li 已提交
3155 3156
  // do config change
  syncNodeDoConfigChange(ths, &newSyncCfg, pEntry->index);
3157

M
Minghao Li 已提交
3158 3159 3160 3161 3162 3163
  // set pFinish
  pFinish->oldCfg = oldSyncCfg;
  pFinish->newCfg = newSyncCfg;
  pFinish->newCfgIndex = pEntry->index;
  pFinish->newCfgTerm = pEntry->term;
  pFinish->newCfgSeqNum = pEntry->seqNum;
3164

M
Minghao Li 已提交
3165 3166
  return 0;
}
3167

M
Minghao Li 已提交
3168 3169 3170 3171 3172 3173 3174 3175 3176
static int32_t syncNodeProposeConfigChangeFinish(SSyncNode* ths, SyncReconfigFinish* pFinish) {
  SRpcMsg rpcMsg;
  syncReconfigFinish2RpcMsg(pFinish, &rpcMsg);

  int32_t code = syncNodePropose(ths, &rpcMsg, false);
  if (code != 0) {
    sError("syncNodeProposeConfigChangeFinish error");
    ths->changing = false;
  }
3177 3178 3179
  return 0;
}

3180 3181 3182 3183
bool syncNodeIsOptimizedOneReplica(SSyncNode* ths, SRpcMsg* pMsg) {
  return (ths->replicaNum == 1 && syncUtilUserCommit(pMsg->msgType) && ths->vgId != 1);
}

M
Minghao Li 已提交
3184
int32_t syncNodeDoCommit(SSyncNode* ths, SyncIndex beginIndex, SyncIndex endIndex, uint64_t flag) {
3185 3186 3187 3188
  if (beginIndex > endIndex) {
    return 0;
  }

M
Minghao Li 已提交
3189 3190 3191 3192 3193 3194 3195 3196 3197 3198 3199 3200 3201
  if (ths == NULL) {
    return -1;
  }

  if (ths->pFsm != NULL && ths->pFsm->FpGetSnapshotInfo != NULL) {
    // advance commit index to sanpshot first
    SSnapshot snapshot = {0};
    ths->pFsm->FpGetSnapshotInfo(ths->pFsm, &snapshot);
    if (snapshot.lastApplyIndex >= 0 && snapshot.lastApplyIndex >= beginIndex) {
      char eventLog[128];
      snprintf(eventLog, sizeof(eventLog), "commit by snapshot from index:%" PRId64 " to index:%" PRId64, beginIndex,
               snapshot.lastApplyIndex);
      syncNodeEventLog(ths, eventLog);
3202

M
Minghao Li 已提交
3203 3204 3205
      // update begin index
      beginIndex = snapshot.lastApplyIndex + 1;
    }
3206 3207
  }

3208 3209
  int32_t    code = 0;
  ESyncState state = flag;
M
Minghao Li 已提交
3210 3211

  char eventLog[128];
3212
  snprintf(eventLog, sizeof(eventLog), "commit by wal from index:%" PRId64 " to index:%" PRId64, beginIndex, endIndex);
M
Minghao Li 已提交
3213
  syncNodeEventLog(ths, eventLog);
3214 3215 3216 3217 3218 3219

  // execute fsm
  if (ths->pFsm != NULL) {
    for (SyncIndex i = beginIndex; i <= endIndex; ++i) {
      if (i != SYNC_INDEX_INVALID) {
        SSyncRaftEntry* pEntry;
3220 3221 3222 3223 3224 3225
        SLRUCache*      pCache = ths->pLogStore->pCache;
        LRUHandle*      h = taosLRUCacheLookup(pCache, &i, sizeof(i));
        if (h) {
          pEntry = (SSyncRaftEntry*)taosLRUCacheValue(pCache, h);
        } else {
          code = ths->pLogStore->syncLogGetEntry(ths->pLogStore, i, &pEntry);
M
Minghao Li 已提交
3226 3227 3228
          // ASSERT(code == 0);
          // ASSERT(pEntry != NULL);
          if (code != 0 || pEntry == NULL) {
M
Minghao Li 已提交
3229
            syncNodeErrorLog(ths, "get log entry error");
3230
            sFatal("vgId:%d, get log entry %" PRId64 " error when commit since %s", ths->vgId, i, terrstr());
M
Minghao Li 已提交
3231 3232
            continue;
          }
3233
        }
3234 3235 3236 3237

        SRpcMsg rpcMsg;
        syncEntry2OriginalRpc(pEntry, &rpcMsg);

3238
        // user commit
3239 3240
        if ((ths->pFsm->FpCommitCb != NULL) && syncUtilUserCommit(pEntry->originalRpcType)) {
          bool internalExecute = true;
S
Shengliang Guan 已提交
3241
          if ((ths->replicaNum == 1) && ths->restoreFinish && ths->vgId != 1) {
3242 3243 3244 3245 3246
            internalExecute = false;
          }

          do {
            char logBuf[128];
S
Shengliang Guan 已提交
3247
            snprintf(logBuf, sizeof(logBuf), "commit index:%" PRId64 ", internal:%d", i, internalExecute);
3248 3249
            syncNodeEventLog(ths, logBuf);
          } while (0);
3250

3251 3252
          // execute fsm in apply thread, or execute outside syncPropose
          if (internalExecute) {
S
Shengliang Guan 已提交
3253 3254 3255 3256 3257 3258 3259 3260 3261 3262 3263 3264 3265 3266
            SFsmCbMeta cbMeta = {
                .index = pEntry->index,
                .lastConfigIndex = syncNodeGetSnapshotConfigIndex(ths, pEntry->index),
                .isWeak = pEntry->isWeak,
                .code = 0,
                .state = ths->state,
                .seqNum = pEntry->seqNum,
                .term = pEntry->term,
                .currentTerm = ths->pRaftStore->currentTerm,
                .flag = flag,
            };

            syncGetAndDelRespRpc(ths, cbMeta.seqNum, &rpcMsg.info);
            ths->pFsm->FpCommitCb(ths->pFsm, &rpcMsg, &cbMeta);
3267
          }
3268 3269 3270
        }

        // config change
3271
        if (pEntry->originalRpcType == TDMT_SYNC_CONFIG_CHANGE) {
M
Minghao Li 已提交
3272 3273 3274 3275 3276 3277 3278 3279 3280 3281 3282 3283 3284 3285
          SyncReconfigFinish* pFinish = syncReconfigFinishBuild(ths->vgId);
          ASSERT(pFinish != NULL);

          code = syncNodeConfigChange(ths, &rpcMsg, pEntry, pFinish);
          ASSERT(code == 0);

          if (ths->state == TAOS_SYNC_STATE_LEADER) {
            syncNodeProposeConfigChangeFinish(ths, pFinish);
          }
          syncReconfigFinishDestroy(pFinish);
        }

        // config change finish
        if (pEntry->originalRpcType == TDMT_SYNC_CONFIG_CHANGE_FINISH) {
M
Minghao Li 已提交
3286
          if (rpcMsg.pCont != NULL && rpcMsg.contLen > 0) {
M
Minghao Li 已提交
3287 3288 3289
            code = syncNodeConfigChangeFinish(ths, &rpcMsg, pEntry);
            ASSERT(code == 0);
          }
3290
        }
3291

3292 3293
#if 0
        // execute in pre-commit
M
Minghao Li 已提交
3294
        // leader transfer
3295 3296 3297
        if (pEntry->originalRpcType == TDMT_SYNC_LEADER_TRANSFER) {
          code = syncDoLeaderTransfer(ths, &rpcMsg, pEntry);
          ASSERT(code == 0);
3298
        }
3299
#endif
3300 3301

        // restore finish
3302
        // if only snapshot, a noop entry will be append, so syncLogLastIndex is always ok
3303 3304 3305 3306 3307 3308
        if (pEntry->index == ths->pLogStore->syncLogLastIndex(ths->pLogStore)) {
          if (ths->restoreFinish == false) {
            if (ths->pFsm->FpRestoreFinishCb != NULL) {
              ths->pFsm->FpRestoreFinishCb(ths->pFsm);
            }
            ths->restoreFinish = true;
M
Minghao Li 已提交
3309

3310 3311
            int64_t restoreDelay = taosGetTimestampMs() - ths->leaderTime;

M
Minghao Li 已提交
3312
            char eventLog[128];
S
Shengliang Guan 已提交
3313 3314
            snprintf(eventLog, sizeof(eventLog), "restore finish, index:%" PRId64 ", elapsed:%" PRId64 " ms, ",
                     pEntry->index, restoreDelay);
M
Minghao Li 已提交
3315
            syncNodeEventLog(ths, eventLog);
3316 3317 3318 3319
          }
        }

        rpcFreeCont(rpcMsg.pCont);
3320 3321 3322 3323 3324
        if (h) {
          taosLRUCacheRelease(pCache, h, false);
        } else {
          syncEntryDestory(pEntry);
        }
3325 3326 3327 3328
      }
    }
  }
  return 0;
3329 3330 3331 3332 3333 3334 3335 3336 3337
}

bool syncNodeInRaftGroup(SSyncNode* ths, SRaftId* pRaftId) {
  for (int i = 0; i < ths->replicaNum; ++i) {
    if (syncUtilSameId(&((ths->replicasId)[i]), pRaftId)) {
      return true;
    }
  }
  return false;
M
Minghao Li 已提交
3338 3339 3340 3341 3342 3343 3344 3345 3346 3347
}

SSyncSnapshotSender* syncNodeGetSnapshotSender(SSyncNode* ths, SRaftId* pDestId) {
  SSyncSnapshotSender* pSender = NULL;
  for (int i = 0; i < ths->replicaNum; ++i) {
    if (syncUtilSameId(pDestId, &((ths->replicasId)[i]))) {
      pSender = (ths->senders)[i];
    }
  }
  return pSender;
M
Minghao Li 已提交
3348
}
M
Minghao Li 已提交
3349

3350 3351 3352 3353 3354 3355 3356 3357 3358 3359
SSyncTimer* syncNodeGetHbTimer(SSyncNode* ths, SRaftId* pDestId) {
  SSyncTimer* pTimer = NULL;
  for (int i = 0; i < ths->replicaNum; ++i) {
    if (syncUtilSameId(pDestId, &((ths->replicasId)[i]))) {
      pTimer = &((ths->peerHeartbeatTimerArr)[i]);
    }
  }
  return pTimer;
}

M
Minghao Li 已提交
3360 3361 3362 3363 3364 3365 3366 3367 3368 3369 3370 3371
SPeerState* syncNodeGetPeerState(SSyncNode* ths, const SRaftId* pDestId) {
  SPeerState* pState = NULL;
  for (int i = 0; i < ths->replicaNum; ++i) {
    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 已提交
3372
  if (pState == NULL) {
3373
    sError("vgId:%d, replica maybe dropped", ths->vgId);
M
Minghao Li 已提交
3374 3375
    return false;
  }
M
Minghao Li 已提交
3376 3377 3378 3379 3380 3381 3382 3383 3384 3385 3386

  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 已提交
3387 3388 3389 3390 3391 3392 3393 3394 3395 3396 3397 3398 3399 3400 3401 3402
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;
    }
  }

  for (int i = 0; i < pSyncNode->peersNum; ++i) {
    SSyncSnapshotSender* pSender = syncNodeGetSnapshotSender(pSyncNode, &(pSyncNode->peersId)[i]);
M
Minghao Li 已提交
3403
    if (pSender != NULL && pSender->start) {
M
Minghao Li 已提交
3404 3405 3406 3407 3408 3409
      sError("sync cannot change3");
      return false;
    }
  }

  return true;
M
Minghao Li 已提交
3410 3411
}

3412 3413 3414 3415 3416 3417 3418 3419 3420 3421 3422 3423 3424 3425
const char* syncTimerTypeStr(enum ESyncTimeoutType timerType) {
  if (timerType == SYNC_TIMEOUT_PING) {
    return "ping";
  } else if (timerType == SYNC_TIMEOUT_ELECTION) {
    return "elect";
  } else if (timerType == SYNC_TIMEOUT_HEARTBEAT) {
    return "heartbeat";
  } else {
    return "unknown";
  }
}

void syncLogRecvTimer(SSyncNode* pSyncNode, const SyncTimeout* pMsg, const char* s) {
  char logBuf[256];
S
Shengliang Guan 已提交
3426
  snprintf(logBuf, sizeof(logBuf), "recv sync-timer {type:%s, lc:%" PRIu64 ", ms:%d, data:%p}, %s",
M
Minghao Li 已提交
3427
           syncTimerTypeStr(pMsg->timeoutType), pMsg->logicClock, pMsg->timerMS, pMsg->data, s);
3428 3429 3430
  syncNodeEventLog(pSyncNode, logBuf);
}

wafwerar's avatar
wafwerar 已提交
3431
void syncLogSendRequestVote(SSyncNode* pSyncNode, const SyncRequestVote* pMsg, const char* s) {
M
Minghao Li 已提交
3432 3433 3434 3435 3436 3437 3438 3439 3440 3441
  char     host[64];
  uint16_t port;
  syncUtilU642Addr(pMsg->destId.addr, host, sizeof(host), &port);
  char logBuf[256];
  snprintf(logBuf, sizeof(logBuf),
           "send sync-request-vote to %s:%d {term:%" PRIu64 ", lindex:%" PRId64 ", lterm:%" PRIu64 "}, %s", host, port,
           pMsg->term, pMsg->lastLogIndex, pMsg->lastLogTerm, s);
  syncNodeEventLog(pSyncNode, logBuf);
}

wafwerar's avatar
wafwerar 已提交
3442
void syncLogRecvRequestVote(SSyncNode* pSyncNode, const SyncRequestVote* pMsg, const char* s) {
M
Minghao Li 已提交
3443 3444 3445 3446 3447 3448 3449 3450 3451 3452
  char     logBuf[256];
  char     host[64];
  uint16_t port;
  syncUtilU642Addr(pMsg->srcId.addr, host, sizeof(host), &port);
  snprintf(logBuf, sizeof(logBuf),
           "recv sync-request-vote from %s:%d, {term:%" PRIu64 ", lindex:%" PRId64 ", lterm:%" PRIu64 "}, %s", host,
           port, pMsg->term, pMsg->lastLogIndex, pMsg->lastLogTerm, s);
  syncNodeEventLog(pSyncNode, logBuf);
}

wafwerar's avatar
wafwerar 已提交
3453
void syncLogSendRequestVoteReply(SSyncNode* pSyncNode, const SyncRequestVoteReply* pMsg, const char* s) {
M
Minghao Li 已提交
3454 3455 3456 3457 3458 3459 3460 3461 3462
  char     host[64];
  uint16_t port;
  syncUtilU642Addr(pMsg->destId.addr, host, sizeof(host), &port);
  char logBuf[256];
  snprintf(logBuf, sizeof(logBuf), "send sync-request-vote-reply to %s:%d {term:%" PRIu64 ", grant:%d}, %s", host, port,
           pMsg->term, pMsg->voteGranted, s);
  syncNodeEventLog(pSyncNode, logBuf);
}

wafwerar's avatar
wafwerar 已提交
3463
void syncLogRecvRequestVoteReply(SSyncNode* pSyncNode, const SyncRequestVoteReply* pMsg, const char* s) {
M
Minghao Li 已提交
3464 3465 3466 3467 3468 3469 3470 3471 3472
  char     host[64];
  uint16_t port;
  syncUtilU642Addr(pMsg->srcId.addr, host, sizeof(host), &port);
  char logBuf[256];
  snprintf(logBuf, sizeof(logBuf), "recv sync-request-vote-reply from %s:%d {term:%" PRIu64 ", grant:%d}, %s", host,
           port, pMsg->term, pMsg->voteGranted, s);
  syncNodeEventLog(pSyncNode, logBuf);
}

wafwerar's avatar
wafwerar 已提交
3473
void syncLogSendAppendEntries(SSyncNode* pSyncNode, const SyncAppendEntries* pMsg, const char* s) {
M
Minghao Li 已提交
3474 3475 3476 3477 3478 3479
  char     host[64];
  uint16_t port;
  syncUtilU642Addr(pMsg->destId.addr, host, sizeof(host), &port);
  char logBuf[256];
  snprintf(logBuf, sizeof(logBuf),
           "send sync-append-entries to %s:%d, {term:%" PRIu64 ", pre-index:%" PRId64 ", pre-term:%" PRIu64
M
Minghao Li 已提交
3480
           ", pterm:%" PRIu64 ", cmt:%" PRId64
M
Minghao Li 已提交
3481 3482 3483 3484 3485 3486 3487
           ", "
           "datalen:%d}, %s",
           host, port, pMsg->term, pMsg->prevLogIndex, pMsg->prevLogTerm, pMsg->privateTerm, pMsg->commitIndex,
           pMsg->dataLen, s);
  syncNodeEventLog(pSyncNode, logBuf);
}

wafwerar's avatar
wafwerar 已提交
3488
void syncLogRecvAppendEntries(SSyncNode* pSyncNode, const SyncAppendEntries* pMsg, const char* s) {
M
Minghao Li 已提交
3489 3490 3491 3492 3493
  char     host[64];
  uint16_t port;
  syncUtilU642Addr(pMsg->srcId.addr, host, sizeof(host), &port);
  char logBuf[256];
  snprintf(logBuf, sizeof(logBuf),
M
Minghao Li 已提交
3494
           "recv sync-append-entries from %s:%d {term:%" PRIu64 ", pre-index:%" PRIu64 ", pre-term:%" PRIu64
M
Minghao Li 已提交
3495
           ", cmt:%" PRIu64 ", pterm:%" PRIu64
M
Minghao Li 已提交
3496 3497 3498 3499
           ", "
           "datalen:%d}, %s",
           host, port, pMsg->term, pMsg->prevLogIndex, pMsg->prevLogTerm, pMsg->commitIndex, pMsg->privateTerm,
           pMsg->dataLen, s);
M
Minghao Li 已提交
3500
  syncNodeEventLog(pSyncNode, logBuf);
M
Minghao Li 已提交
3501 3502
}

wafwerar's avatar
wafwerar 已提交
3503
void syncLogSendAppendEntriesBatch(SSyncNode* pSyncNode, const SyncAppendEntriesBatch* pMsg, const char* s) {
M
Minghao Li 已提交
3504 3505 3506 3507 3508 3509
  char     host[64];
  uint16_t port;
  syncUtilU642Addr(pMsg->destId.addr, host, sizeof(host), &port);
  char logBuf[256];
  snprintf(logBuf, sizeof(logBuf),
           "send sync-append-entries-batch to %s:%d, {term:%" PRIu64 ", pre-index:%" PRId64 ", pre-term:%" PRIu64
M
Minghao Li 已提交
3510
           ", pterm:%" PRIu64 ", cmt:%" PRId64 ", datalen:%d, count:%d}, %s",
M
Minghao Li 已提交
3511 3512
           host, port, pMsg->term, pMsg->prevLogIndex, pMsg->prevLogTerm, pMsg->privateTerm, pMsg->commitIndex,
           pMsg->dataLen, pMsg->dataCount, s);
M
Minghao Li 已提交
3513 3514 3515
  syncNodeEventLog(pSyncNode, logBuf);
}

wafwerar's avatar
wafwerar 已提交
3516
void syncLogRecvAppendEntriesBatch(SSyncNode* pSyncNode, const SyncAppendEntriesBatch* pMsg, const char* s) {
M
Minghao Li 已提交
3517 3518 3519 3520 3521 3522
  char     host[64];
  uint16_t port;
  syncUtilU642Addr(pMsg->srcId.addr, host, sizeof(host), &port);
  char logBuf[256];
  snprintf(logBuf, sizeof(logBuf),
           "recv sync-append-entries-batch from %s:%d, {term:%" PRIu64 ", pre-index:%" PRId64 ", pre-term:%" PRIu64
M
Minghao Li 已提交
3523
           ", pterm:%" PRIu64 ", cmt:%" PRId64 ", datalen:%d, count:%d}, %s",
M
Minghao Li 已提交
3524 3525
           host, port, pMsg->term, pMsg->prevLogIndex, pMsg->prevLogTerm, pMsg->privateTerm, pMsg->commitIndex,
           pMsg->dataLen, pMsg->dataCount, s);
M
Minghao Li 已提交
3526
  syncNodeEventLog(pSyncNode, logBuf);
M
Minghao Li 已提交
3527 3528
}

3529
void syncLogSendAppendEntriesReply(SSyncNode* pSyncNode, const SyncAppendEntriesReply* pMsg, const char* s) {
M
Minghao Li 已提交
3530 3531 3532 3533 3534 3535 3536 3537 3538 3539 3540
  char     host[64];
  uint16_t port;
  syncUtilU642Addr(pMsg->destId.addr, host, sizeof(host), &port);
  char logBuf[256];
  snprintf(logBuf, sizeof(logBuf),
           "send sync-append-entries-reply to %s:%d, {term:%" PRIu64 ", pterm:%" PRIu64 ", success:%d, match:%" PRId64
           "}, %s",
           host, port, pMsg->term, pMsg->privateTerm, pMsg->success, pMsg->matchIndex, s);
  syncNodeEventLog(pSyncNode, logBuf);
}

wafwerar's avatar
wafwerar 已提交
3541
void syncLogRecvAppendEntriesReply(SSyncNode* pSyncNode, const SyncAppendEntriesReply* pMsg, const char* s) {
M
Minghao Li 已提交
3542 3543 3544 3545 3546
  char     host[64];
  uint16_t port;
  syncUtilU642Addr(pMsg->srcId.addr, host, sizeof(host), &port);
  char logBuf[256];
  snprintf(logBuf, sizeof(logBuf),
M
Minghao Li 已提交
3547
           "recv sync-append-entries-reply from %s:%d {term:%" PRIu64 ", pterm:%" PRIu64 ", success:%d, match:%" PRId64
M
Minghao Li 已提交
3548 3549
           "}, %s",
           host, port, pMsg->term, pMsg->privateTerm, pMsg->success, pMsg->matchIndex, s);
M
Minghao Li 已提交
3550
  syncNodeEventLog(pSyncNode, logBuf);
M
Minghao Li 已提交
3551
}
3552 3553 3554 3555 3556 3557 3558

void syncLogSendHeartbeat(SSyncNode* pSyncNode, const SyncHeartbeat* pMsg, const char* s) {
  char     host[64];
  uint16_t port;
  syncUtilU642Addr(pMsg->destId.addr, host, sizeof(host), &port);
  char logBuf[256];
  snprintf(logBuf, sizeof(logBuf),
M
Minghao Li 已提交
3559
           "send sync-heartbeat to %s:%d {term:%" PRIu64 ", cmt:%" PRId64 ", min-match:%" PRId64 ", pterm:%" PRIu64
M
Minghao Li 已提交
3560 3561
           "}, %s",
           host, port, pMsg->term, pMsg->commitIndex, pMsg->minMatchIndex, pMsg->privateTerm, s);
3562 3563 3564 3565 3566 3567 3568 3569 3570
  syncNodeEventLog(pSyncNode, logBuf);
}

void syncLogRecvHeartbeat(SSyncNode* pSyncNode, const SyncHeartbeat* pMsg, const char* s) {
  char     host[64];
  uint16_t port;
  syncUtilU642Addr(pMsg->srcId.addr, host, sizeof(host), &port);
  char logBuf[256];
  snprintf(logBuf, sizeof(logBuf),
M
Minghao Li 已提交
3571 3572 3573
           "recv sync-heartbeat from %s:%d {term:%" PRIu64 ", cmt:%" PRId64 ", min-match:%" PRId64 ", pterm:%" PRIu64
           "}, %s",
           host, port, pMsg->term, pMsg->commitIndex, pMsg->minMatchIndex, pMsg->privateTerm, s);
3574 3575 3576 3577 3578 3579 3580 3581 3582 3583 3584 3585 3586 3587 3588 3589 3590 3591 3592 3593 3594
  syncNodeEventLog(pSyncNode, logBuf);
}

void syncLogSendHeartbeatReply(SSyncNode* pSyncNode, const SyncHeartbeatReply* pMsg, const char* s) {
  char     host[64];
  uint16_t port;
  syncUtilU642Addr(pMsg->destId.addr, host, sizeof(host), &port);
  char logBuf[256];
  snprintf(logBuf, sizeof(logBuf), "send sync-heartbeat-reply from %s:%d {term:%" PRIu64 ", pterm:%" PRIu64 "}, %s",
           host, port, pMsg->term, pMsg->privateTerm, s);
  syncNodeEventLog(pSyncNode, logBuf);
}

void syncLogRecvHeartbeatReply(SSyncNode* pSyncNode, const SyncHeartbeatReply* pMsg, const char* s) {
  char     host[64];
  uint16_t port;
  syncUtilU642Addr(pMsg->srcId.addr, host, sizeof(host), &port);
  char logBuf[256];
  snprintf(logBuf, sizeof(logBuf), "recv sync-heartbeat-reply from %s:%d {term:%" PRIu64 ", pterm:%" PRIu64 "}, %s",
           host, port, pMsg->term, pMsg->privateTerm, s);
  syncNodeEventLog(pSyncNode, logBuf);
M
Minghao Li 已提交
3595
}
3596 3597 3598 3599 3600 3601

void syncLogRecvLocalCmd(SSyncNode* pSyncNode, const SyncLocalCmd* pMsg, const char* s) {
  char logBuf[256];
  snprintf(logBuf, sizeof(logBuf), "recv sync-local-cmd {cmd:%d-%s, sd-new-term:%" PRIu64 "}, %s", pMsg->cmd,
           syncLocalCmdGetStr(pMsg->cmd), pMsg->sdNewTerm, s);
  syncNodeEventLog(pSyncNode, logBuf);
S
Shengliang Guan 已提交
3602
}