syncMain.c 75.2 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/>.
 */

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

38
bool gRaftDetailLog = false;
39

M
Minghao Li 已提交
40 41 42
static int32_t tsNodeRefId = -1;

// ------ local funciton ---------
M
Minghao Li 已提交
43
// enqueue message ----
M
Minghao Li 已提交
44 45 46 47 48
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);
M
Minghao Li 已提交
49

M
Minghao Li 已提交
50
// process message ----
M
Minghao Li 已提交
51 52 53
int32_t syncNodeOnPingCb(SSyncNode* ths, SyncPing* pMsg);
int32_t syncNodeOnPingReplyCb(SSyncNode* ths, SyncPingReply* pMsg);
int32_t syncNodeOnClientRequestCb(SSyncNode* ths, SyncClientRequest* pMsg);
M
Minghao Li 已提交
54 55 56

// life cycle
static void syncFreeNode(void* param);
M
Minghao Li 已提交
57 58 59
// ---------------------------------

int32_t syncInit() {
M
Minghao Li 已提交
60 61 62 63 64 65 66 67 68 69 70
  int32_t ret = 0;

  if (!syncEnvIsStart()) {
    tsNodeRefId = taosOpenRef(200, syncFreeNode);
    if (tsNodeRefId < 0) {
      sError("failed to init node ref");
      syncCleanUp();
      ret = -1;
    } else {
      ret = syncEnvStart();
    }
M
Minghao Li 已提交
71 72
  }

M
Minghao Li 已提交
73
  return ret;
M
Minghao Li 已提交
74
}
M
Minghao Li 已提交
75

M
Minghao Li 已提交
76 77 78
void syncCleanUp() {
  int32_t ret = syncEnvStop();
  assert(ret == 0);
M
Minghao Li 已提交
79 80 81 82 83

  if (tsNodeRefId != -1) {
    taosCloseRef(tsNodeRefId);
    tsNodeRefId = -1;
  }
M
Minghao Li 已提交
84
}
M
Minghao Li 已提交
85

M
Minghao Li 已提交
86
int64_t syncOpen(const SSyncInfo* pSyncInfo) {
M
Minghao Li 已提交
87 88
  SSyncNode* pSyncNode = syncNodeOpen(pSyncInfo);
  assert(pSyncNode != NULL);
M
Minghao Li 已提交
89

90 91 92
  if (gRaftDetailLog) {
    syncNodeLog2("syncNodeOpen open success", pSyncNode);
  }
M
Minghao Li 已提交
93

M
Minghao Li 已提交
94 95 96 97 98 99 100
  pSyncNode->rid = taosAddRef(tsNodeRefId, pSyncNode);
  if (pSyncNode->rid < 0) {
    syncFreeNode(pSyncNode);
    return -1;
  }

  return pSyncNode->rid;
M
Minghao Li 已提交
101
}
M
Minghao Li 已提交
102

M
Minghao Li 已提交
103 104 105 106 107
void syncStart(int64_t rid) {
  SSyncNode* pSyncNode = (SSyncNode*)taosAcquireRef(tsNodeRefId, rid);
  if (pSyncNode == NULL) {
    return;
  }
M
Minghao Li 已提交
108 109 110 111 112 113 114 115 116 117 118 119 120 121 122

  if (pSyncNode->pRaftCfg->isStandBy) {
    syncNodeStartStandBy(pSyncNode);
  } else {
    syncNodeStart(pSyncNode);
  }

  taosReleaseRef(tsNodeRefId, pSyncNode->rid);
}

void syncStartNormal(int64_t rid) {
  SSyncNode* pSyncNode = (SSyncNode*)taosAcquireRef(tsNodeRefId, rid);
  if (pSyncNode == NULL) {
    return;
  }
M
Minghao Li 已提交
123 124 125 126 127
  syncNodeStart(pSyncNode);

  taosReleaseRef(tsNodeRefId, pSyncNode->rid);
}

M
Minghao Li 已提交
128 129 130 131 132 133 134 135 136 137
void syncStartStandBy(int64_t rid) {
  SSyncNode* pSyncNode = (SSyncNode*)taosAcquireRef(tsNodeRefId, rid);
  if (pSyncNode == NULL) {
    return;
  }
  syncNodeStartStandBy(pSyncNode);

  taosReleaseRef(tsNodeRefId, pSyncNode->rid);
}

M
Minghao Li 已提交
138
void syncStop(int64_t rid) {
M
Minghao Li 已提交
139 140 141 142
  SSyncNode* pSyncNode = (SSyncNode*)taosAcquireRef(tsNodeRefId, rid);
  if (pSyncNode == NULL) {
    return;
  }
M
Minghao Li 已提交
143
  syncNodeClose(pSyncNode);
M
Minghao Li 已提交
144 145 146

  taosReleaseRef(tsNodeRefId, pSyncNode->rid);
  taosRemoveRef(tsNodeRefId, rid);
M
Minghao Li 已提交
147
}
M
Minghao Li 已提交
148

M
Minghao Li 已提交
149 150 151
int32_t syncSetStandby(int64_t rid) {
  SSyncNode* pSyncNode = (SSyncNode*)taosAcquireRef(tsNodeRefId, rid);
  if (pSyncNode == NULL) {
M
Minghao Li 已提交
152 153
    terrno = TSDB_CODE_SYN_INTERNAL_ERROR;
    return -1;
M
Minghao Li 已提交
154 155
  }

156
  if (pSyncNode->state == TAOS_SYNC_STATE_LEADER) {
M
Minghao Li 已提交
157
    taosReleaseRef(tsNodeRefId, pSyncNode->rid);
M
Minghao Li 已提交
158 159
    terrno = TSDB_CODE_SYN_INTERNAL_ERROR;
    return -1;
M
Minghao Li 已提交
160 161 162 163 164 165 166 167 168 169 170 171 172 173 174
  }

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

  pSyncNode->pRaftCfg->isStandBy = 1;
  raftCfgPersist(pSyncNode->pRaftCfg);

  taosReleaseRef(tsNodeRefId, pSyncNode->rid);
175
  sInfo("vgId:%d, set to standby", pSyncNode->vgId);
M
Minghao Li 已提交
176 177 178
  return 0;
}

M
Minghao Li 已提交
179 180 181
int32_t syncReconfigBuild(int64_t rid, const SSyncCfg* pNewCfg, SRpcMsg* pRpcMsg) {
  SSyncNode* pSyncNode = (SSyncNode*)taosAcquireRef(tsNodeRefId, rid);
  if (pSyncNode == NULL) {
M
Minghao Li 已提交
182 183
    terrno = TSDB_CODE_SYN_INTERNAL_ERROR;
    return -1;
M
Minghao Li 已提交
184 185 186
  }
  ASSERT(rid == pSyncNode->rid);

M
Minghao Li 已提交
187
  int32_t ret = 0;
188 189 190
  bool    IamInNew = syncNodeInConfig(pSyncNode, pNewCfg);

#if 0
M
Minghao Li 已提交
191
  for (int i = 0; i < pNewCfg->replicaNum; ++i) {
M
Minghao Li 已提交
192 193
    if (strcmp((pNewCfg->nodeInfo)[i].nodeFqdn, pSyncNode->myNodeInfo.nodeFqdn) == 0 &&
        (pNewCfg->nodeInfo)[i].nodePort == pSyncNode->myNodeInfo.nodePort) {
M
Minghao Li 已提交
194 195
      IamInNew = true;
    }
M
Minghao Li 已提交
196 197 198 199 200 201 202 203 204

    /*
        SRaftId newId;
        newId.addr = syncUtilAddr2U64((pNewCfg->nodeInfo)[i].nodeFqdn, (pNewCfg->nodeInfo)[i].nodePort);
        newId.vgId = pSyncNode->vgId;
        if (syncUtilSameId(&(pSyncNode->myRaftId), &newId)) {
          IamInNew = true;
        }
    */
M
Minghao Li 已提交
205
  }
206
#endif
M
Minghao Li 已提交
207

M
Minghao Li 已提交
208 209
  if (!IamInNew) {
    taosReleaseRef(tsNodeRefId, pSyncNode->rid);
M
Minghao Li 已提交
210 211
    terrno = TSDB_CODE_SYN_NOT_IN_NEW_CONFIG;
    return -1;
M
Minghao Li 已提交
212 213 214 215 216 217 218 219 220
  }

  char* newconfig = syncCfg2Str((SSyncCfg*)pNewCfg);
  pRpcMsg->msgType = TDMT_SYNC_CONFIG_CHANGE;
  pRpcMsg->info.noResp = 1;
  pRpcMsg->contLen = strlen(newconfig) + 1;
  pRpcMsg->pCont = rpcMallocCont(pRpcMsg->contLen);
  snprintf(pRpcMsg->pCont, pRpcMsg->contLen, "%s", newconfig);
  taosMemoryFree(newconfig);
221

M
Minghao Li 已提交
222 223 224 225 226 227 228
  taosReleaseRef(tsNodeRefId, pSyncNode->rid);
  return ret;
}

int32_t syncReconfig(int64_t rid, const SSyncCfg* pNewCfg) {
  SSyncNode* pSyncNode = (SSyncNode*)taosAcquireRef(tsNodeRefId, rid);
  if (pSyncNode == NULL) {
M
Minghao Li 已提交
229 230
    terrno = TSDB_CODE_SYN_INTERNAL_ERROR;
    return -1;
M
Minghao Li 已提交
231 232 233
  }
  ASSERT(rid == pSyncNode->rid);

234 235 236
  bool IamInNew = syncNodeInConfig(pSyncNode, pNewCfg);

#if 0 
M
Minghao Li 已提交
237
  for (int i = 0; i < pNewCfg->replicaNum; ++i) {
238 239
    if (strcmp((pNewCfg->nodeInfo)[i].nodeFqdn, pSyncNode->myNodeInfo.nodeFqdn) == 0 &&
        (pNewCfg->nodeInfo)[i].nodePort == pSyncNode->myNodeInfo.nodePort) {
M
Minghao Li 已提交
240 241
      IamInNew = true;
    }
242 243 244 245 246 247 248 249 250 251 252 253

    /*
        // some problem in inet_addr

        SRaftId newId = EMPTY_RAFT_ID;
        newId.addr = syncUtilAddr2U64((pNewCfg->nodeInfo)[i].nodeFqdn, (pNewCfg->nodeInfo)[i].nodePort);
        newId.vgId = pSyncNode->vgId;

        if (syncUtilSameId(&(pSyncNode->myRaftId), &newId)) {
          IamInNew = true;
        }
      */
M
Minghao Li 已提交
254
  }
255
#endif
M
Minghao Li 已提交
256

M
Minghao Li 已提交
257
  if (!IamInNew) {
M
Minghao Li 已提交
258
    sError("sync reconfig error, not in new config");
M
Minghao Li 已提交
259
    taosReleaseRef(tsNodeRefId, pSyncNode->rid);
M
Minghao Li 已提交
260 261
    terrno = TSDB_CODE_SYN_NOT_IN_NEW_CONFIG;
    return -1;
M
Minghao Li 已提交
262
  }
263

M
Minghao Li 已提交
264
  char* newconfig = syncCfg2Str((SSyncCfg*)pNewCfg);
265
  if (gRaftDetailLog) {
266
    sInfo("==syncReconfig== newconfig:%s", newconfig);
267
  }
M
Minghao Li 已提交
268

M
Minghao Li 已提交
269 270
  int32_t ret = 0;

M
Minghao Li 已提交
271
  SRpcMsg rpcMsg = {0};
272
  rpcMsg.msgType = TDMT_SYNC_CONFIG_CHANGE;
S
Shengliang Guan 已提交
273
  rpcMsg.info.noResp = 1;
274
  rpcMsg.contLen = strlen(newconfig) + 1;
M
Minghao Li 已提交
275
  rpcMsg.pCont = rpcMallocCont(rpcMsg.contLen);
276 277
  snprintf(rpcMsg.pCont, rpcMsg.contLen, "%s", newconfig);
  taosMemoryFree(newconfig);
M
Minghao Li 已提交
278 279 280
  ret = syncNodePropose(pSyncNode, &rpcMsg, false);

  taosReleaseRef(tsNodeRefId, pSyncNode->rid);
M
Minghao Li 已提交
281 282
  return ret;
}
M
Minghao Li 已提交
283

284
int32_t syncLeaderTransfer(int64_t rid) {
M
Minghao Li 已提交
285 286
  SSyncNode* pSyncNode = (SSyncNode*)taosAcquireRef(tsNodeRefId, rid);
  if (pSyncNode == NULL) {
M
Minghao Li 已提交
287 288
    terrno = TSDB_CODE_SYN_INTERNAL_ERROR;
    return -1;
M
Minghao Li 已提交
289 290 291
  }
  ASSERT(rid == pSyncNode->rid);

M
Minghao Li 已提交
292
  if (pSyncNode->peersNum == 0) {
M
Minghao Li 已提交
293
    taosReleaseRef(tsNodeRefId, pSyncNode->rid);
M
Minghao Li 已提交
294 295
    terrno = TSDB_CODE_SYN_INTERNAL_ERROR;
    return -1;
M
Minghao Li 已提交
296 297 298 299
  }

  SNodeInfo newLeader = (pSyncNode->peersNodeInfo)[0];
  taosReleaseRef(tsNodeRefId, pSyncNode->rid);
300

M
Minghao Li 已提交
301
  int32_t ret = syncLeaderTransferTo(rid, newLeader);
302 303 304 305 306 307
  return ret;
}

int32_t syncLeaderTransferTo(int64_t rid, SNodeInfo newLeader) {
  SSyncNode* pSyncNode = (SSyncNode*)taosAcquireRef(tsNodeRefId, rid);
  if (pSyncNode == NULL) {
M
Minghao Li 已提交
308 309
    terrno = TSDB_CODE_SYN_INTERNAL_ERROR;
    return -1;
310
  }
M
Minghao Li 已提交
311
  ASSERT(rid == pSyncNode->rid);
312 313 314 315
  int32_t ret = 0;

  if (pSyncNode->replicaNum == 1) {
    sError("only one replica, cannot drop leader");
M
Minghao Li 已提交
316
    taosReleaseRef(tsNodeRefId, pSyncNode->rid);
M
Minghao Li 已提交
317 318
    terrno = TSDB_CODE_SYN_ONE_REPLICA;
    return -1;
319 320 321 322 323
  }

  SyncLeaderTransfer* pMsg = syncLeaderTransferBuild(pSyncNode->vgId);
  pMsg->newLeaderId.addr = syncUtilAddr2U64(newLeader.nodeFqdn, newLeader.nodePort);
  pMsg->newLeaderId.vgId = pSyncNode->vgId;
M
Minghao Li 已提交
324
  pMsg->newNodeInfo = newLeader;
325 326 327 328 329
  ASSERT(pMsg != NULL);
  SRpcMsg rpcMsg = {0};
  syncLeaderTransfer2RpcMsg(pMsg, &rpcMsg);
  syncLeaderTransferDestroy(pMsg);

M
Minghao Li 已提交
330
  ret = syncNodePropose(pSyncNode, &rpcMsg, false);
331 332 333 334 335 336 337 338 339 340 341 342 343 344 345 346 347 348 349 350 351 352 353 354 355 356 357 358 359 360 361 362 363 364 365 366
  taosReleaseRef(tsNodeRefId, pSyncNode->rid);
  return ret;
}

bool syncCanLeaderTransfer(int64_t rid) {
  SSyncNode* pSyncNode = (SSyncNode*)taosAcquireRef(tsNodeRefId, rid);
  if (pSyncNode == NULL) {
    return false;
  }
  assert(rid == pSyncNode->rid);

  if (pSyncNode->replicaNum == 1) {
    taosReleaseRef(tsNodeRefId, pSyncNode->rid);
    return false;
  }

  if (pSyncNode->state == TAOS_SYNC_STATE_FOLLOWER) {
    taosReleaseRef(tsNodeRefId, pSyncNode->rid);
    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;
      }
    }
  }

  taosReleaseRef(tsNodeRefId, pSyncNode->rid);
  return matchOK;
}

M
Minghao Li 已提交
367 368 369 370
int32_t syncForwardToPeer(int64_t rid, const SRpcMsg* pMsg, bool isWeak) {
  int32_t ret = syncPropose(rid, pMsg, isWeak);
  return ret;
}
M
Minghao Li 已提交
371

M
Minghao Li 已提交
372 373 374 375 376 377
ESyncState syncGetMyRole(int64_t rid) {
  SSyncNode* pSyncNode = (SSyncNode*)taosAcquireRef(tsNodeRefId, rid);
  if (pSyncNode == NULL) {
    return TAOS_SYNC_STATE_ERROR;
  }
  assert(rid == pSyncNode->rid);
M
Minghao Li 已提交
378 379 380 381
  ESyncState state = pSyncNode->state;

  taosReleaseRef(tsNodeRefId, pSyncNode->rid);
  return state;
M
Minghao Li 已提交
382 383
}

M
Minghao Li 已提交
384 385 386 387 388 389 390 391 392 393 394 395
bool syncIsRestoreFinish(int64_t rid) {
  SSyncNode* pSyncNode = (SSyncNode*)taosAcquireRef(tsNodeRefId, rid);
  if (pSyncNode == NULL) {
    return false;
  }
  assert(rid == pSyncNode->rid);
  bool b = pSyncNode->restoreFinish;

  taosReleaseRef(tsNodeRefId, pSyncNode->rid);
  return b;
}

396 397 398 399 400 401
int32_t syncGetSnapshotMeta(int64_t rid, struct SSnapshotMeta* sMeta) {
  SSyncNode* pSyncNode = (SSyncNode*)taosAcquireRef(tsNodeRefId, rid);
  if (pSyncNode == NULL) {
    return -1;
  }
  assert(rid == pSyncNode->rid);
402 403 404
  sMeta->lastConfigIndex = pSyncNode->pRaftCfg->lastConfigIndex;

  sTrace("sync get snapshot meta: lastConfigIndex:%ld", pSyncNode->pRaftCfg->lastConfigIndex);
405 406 407 408 409

  taosReleaseRef(tsNodeRefId, pSyncNode->rid);
  return 0;
}

M
Minghao Li 已提交
410 411 412 413 414
const char* syncGetMyRoleStr(int64_t rid) {
  const char* s = syncUtilState2String(syncGetMyRole(rid));
  return s;
}

M
Minghao Li 已提交
415 416 417 418 419 420 421 422 423 424 425 426
int32_t syncGetVgId(int64_t rid) {
  SSyncNode* pSyncNode = (SSyncNode*)taosAcquireRef(tsNodeRefId, rid);
  if (pSyncNode == NULL) {
    return TAOS_SYNC_STATE_ERROR;
  }
  assert(rid == pSyncNode->rid);
  int32_t vgId = pSyncNode->vgId;

  taosReleaseRef(tsNodeRefId, pSyncNode->rid);
  return vgId;
}

M
Minghao Li 已提交
427
SyncTerm syncGetMyTerm(int64_t rid) {
M
Minghao Li 已提交
428 429
  SSyncNode* pSyncNode = (SSyncNode*)taosAcquireRef(tsNodeRefId, rid);
  if (pSyncNode == NULL) {
M
Minghao Li 已提交
430 431 432 433 434 435 436 437 438
    return TAOS_SYNC_STATE_ERROR;
  }
  assert(rid == pSyncNode->rid);
  SyncTerm term = pSyncNode->pRaftStore->currentTerm;

  taosReleaseRef(tsNodeRefId, pSyncNode->rid);
  return term;
}

M
Minghao Li 已提交
439 440 441 442 443 444 445 446 447
void syncGetEpSet(int64_t rid, SEpSet* pEpSet) {
  SSyncNode* pSyncNode = (SSyncNode*)taosAcquireRef(tsNodeRefId, rid);
  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) {
M
Minghao Li 已提交
448 449
    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 已提交
450
    (pEpSet->numOfEps)++;
M
Minghao Li 已提交
451 452

    sInfo("syncGetEpSet index:%d %s:%d", i, pEpSet->eps[i].fqdn, pEpSet->eps[i].port);
M
Minghao Li 已提交
453 454
  }
  pEpSet->inUse = pSyncNode->pRaftCfg->cfg.myIndex;
M
Minghao Li 已提交
455

M
Minghao Li 已提交
456
  sInfo("syncGetEpSet pEpSet->inUse:%d ", pEpSet->inUse);
M
Minghao Li 已提交
457 458 459 460

  taosReleaseRef(tsNodeRefId, pSyncNode->rid);
}

M
Minghao Li 已提交
461 462 463 464 465 466 467 468 469 470 471 472 473 474 475 476 477
int32_t syncGetRespRpc(int64_t rid, uint64_t index, SRpcMsg* msg) {
  SSyncNode* pSyncNode = (SSyncNode*)taosAcquireRef(tsNodeRefId, rid);
  if (pSyncNode == NULL) {
    return TAOS_SYNC_STATE_ERROR;
  }
  assert(rid == pSyncNode->rid);

  SRespStub stub;
  int32_t   ret = syncRespMgrGet(pSyncNode->pSyncRespMgr, index, &stub);
  if (ret == 1) {
    memcpy(msg, &(stub.rpcMsg), sizeof(SRpcMsg));
  }

  taosReleaseRef(tsNodeRefId, pSyncNode->rid);
  return ret;
}

S
Shengliang Guan 已提交
478
int32_t syncGetAndDelRespRpc(int64_t rid, uint64_t index, SRpcHandleInfo* pInfo) {
M
Minghao Li 已提交
479 480 481 482 483 484 485 486 487
  SSyncNode* pSyncNode = (SSyncNode*)taosAcquireRef(tsNodeRefId, rid);
  if (pSyncNode == NULL) {
    return TAOS_SYNC_STATE_ERROR;
  }
  assert(rid == pSyncNode->rid);

  SRespStub stub;
  int32_t   ret = syncRespMgrGetAndDel(pSyncNode->pSyncRespMgr, index, &stub);
  if (ret == 1) {
S
Shengliang Guan 已提交
488
    *pInfo = stub.rpcMsg.info;
M
Minghao Li 已提交
489 490
  }

S
Shengliang Guan 已提交
491
  sTrace("vgId:%d, get seq:%" PRIu64 " rpc handle:%p", pSyncNode->vgId, index, pInfo->handle);
M
Minghao Li 已提交
492 493 494 495
  taosReleaseRef(tsNodeRefId, pSyncNode->rid);
  return ret;
}

496
void syncSetMsgCb(int64_t rid, const SMsgCb* msgcb) {
M
Minghao Li 已提交
497 498 499 500 501 502
  SSyncNode* pSyncNode = (SSyncNode*)taosAcquireRef(tsNodeRefId, rid);
  if (pSyncNode == NULL) {
    sTrace("syncSetQ get pSyncNode is NULL, rid:%ld", rid);
    return;
  }
  assert(rid == pSyncNode->rid);
S
Shengliang Guan 已提交
503
  pSyncNode->msgcb = msgcb;
M
Minghao Li 已提交
504 505 506 507 508 509 510 511 512 513 514 515 516 517 518 519 520 521 522 523 524 525 526 527 528 529 530 531 532 533 534 535 536 537 538 539 540 541 542 543 544 545 546 547 548 549 550 551 552 553 554 555 556

  taosReleaseRef(tsNodeRefId, pSyncNode->rid);
}

char* sync2SimpleStr(int64_t rid) {
  SSyncNode* pSyncNode = (SSyncNode*)taosAcquireRef(tsNodeRefId, rid);
  if (pSyncNode == NULL) {
    sTrace("syncSetRpc get pSyncNode is NULL, rid:%ld", rid);
    return NULL;
  }
  assert(rid == pSyncNode->rid);
  char* s = syncNode2SimpleStr(pSyncNode);
  taosReleaseRef(tsNodeRefId, pSyncNode->rid);

  return s;
}

void setPingTimerMS(int64_t rid, int32_t pingTimerMS) {
  SSyncNode* pSyncNode = (SSyncNode*)taosAcquireRef(tsNodeRefId, rid);
  if (pSyncNode == NULL) {
    return;
  }
  assert(rid == pSyncNode->rid);
  pSyncNode->pingBaseLine = pingTimerMS;
  pSyncNode->pingTimerMS = pingTimerMS;

  taosReleaseRef(tsNodeRefId, pSyncNode->rid);
}

void setElectTimerMS(int64_t rid, int32_t electTimerMS) {
  SSyncNode* pSyncNode = (SSyncNode*)taosAcquireRef(tsNodeRefId, rid);
  if (pSyncNode == NULL) {
    return;
  }
  assert(rid == pSyncNode->rid);
  pSyncNode->electBaseLine = electTimerMS;

  taosReleaseRef(tsNodeRefId, pSyncNode->rid);
}

void setHeartbeatTimerMS(int64_t rid, int32_t hbTimerMS) {
  SSyncNode* pSyncNode = (SSyncNode*)taosAcquireRef(tsNodeRefId, rid);
  if (pSyncNode == NULL) {
    return;
  }
  assert(rid == pSyncNode->rid);
  pSyncNode->hbBaseLine = hbTimerMS;
  pSyncNode->heartbeatTimerMS = hbTimerMS;

  taosReleaseRef(tsNodeRefId, pSyncNode->rid);
}

int32_t syncPropose(int64_t rid, const SRpcMsg* pMsg, bool isWeak) {
M
Minghao Li 已提交
557
  int32_t ret = 0;
M
Minghao Li 已提交
558

559
  SSyncNode* pSyncNode = taosAcquireRef(tsNodeRefId, rid);
560
  if (pSyncNode == NULL) {
M
Minghao Li 已提交
561 562
    terrno = TSDB_CODE_SYN_INTERNAL_ERROR;
    return -1;
563
  }
M
Minghao Li 已提交
564
  assert(rid == pSyncNode->rid);
565 566
  sDebug("vgId:%d sync event currentTerm:%lu propose msgType:%s,%d", pSyncNode->vgId,
         pSyncNode->pRaftStore->currentTerm, TMSG_INFO(pMsg->msgType), pMsg->msgType);
M
Minghao Li 已提交
567 568 569 570 571 572 573
  ret = syncNodePropose(pSyncNode, pMsg, isWeak);

  taosReleaseRef(tsNodeRefId, pSyncNode->rid);
  return ret;
}

int32_t syncNodePropose(SSyncNode* pSyncNode, const SRpcMsg* pMsg, bool isWeak) {
M
Minghao Li 已提交
574
  int32_t ret = 0;
575 576
  sDebug("vgId:%d sync event currentTerm:%lu propose msgType:%s,%d", pSyncNode->vgId,
         pSyncNode->pRaftStore->currentTerm, TMSG_INFO(pMsg->msgType), pMsg->msgType);
M
Minghao Li 已提交
577

M
Minghao Li 已提交
578
  if (pSyncNode->state == TAOS_SYNC_STATE_LEADER) {
M
Minghao Li 已提交
579 580 581 582 583 584
    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 已提交
585 586
    SRpcMsg            rpcMsg;
    syncClientRequest2RpcMsg(pSyncMsg, &rpcMsg);
587 588

    if (pSyncNode->FpEqMsg != NULL && (*pSyncNode->FpEqMsg)(pSyncNode->msgcb, &rpcMsg) == 0) {
M
Minghao Li 已提交
589
      ret = 0;
M
Minghao Li 已提交
590
    } else {
M
Minghao Li 已提交
591 592
      ret = -1;
      terrno = TSDB_CODE_SYN_INTERNAL_ERROR;
M
Minghao Li 已提交
593
      sError("syncPropose pSyncNode->FpEqMsg is NULL");
M
Minghao Li 已提交
594
    }
M
Minghao Li 已提交
595 596
    syncClientRequestDestroy(pSyncMsg);
  } else {
M
Minghao Li 已提交
597 598
    ret = -1;
    terrno = TSDB_CODE_SYN_NOT_LEADER;
M
Minghao Li 已提交
599
    sError("syncPropose not leader, %s", syncUtilState2String(pSyncNode->state));
M
Minghao Li 已提交
600
  }
M
Minghao Li 已提交
601 602 603 604

  return ret;
}

M
Minghao Li 已提交
605
// open/close --------------
606 607 608
SSyncNode* syncNodeOpen(const SSyncInfo* pOldSyncInfo) {
  SSyncInfo* pSyncInfo = (SSyncInfo*)pOldSyncInfo;

wafwerar's avatar
wafwerar 已提交
609
  SSyncNode* pSyncNode = (SSyncNode*)taosMemoryMalloc(sizeof(SSyncNode));
M
Minghao Li 已提交
610
  assert(pSyncNode != NULL);
M
Minghao Li 已提交
611
  memset(pSyncNode, 0, sizeof(SSyncNode));
M
Minghao Li 已提交
612

M
Minghao Li 已提交
613 614 615 616 617 618 619
  int32_t ret = 0;
  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());
      return NULL;
    }
620
  }
M
Minghao Li 已提交
621

622 623
  snprintf(pSyncNode->configPath, sizeof(pSyncNode->configPath), "%s/raft_config.json", pSyncInfo->path);
  if (!taosCheckExistFile(pSyncNode->configPath)) {
M
Minghao Li 已提交
624 625 626 627
    // create a new raft config file
    SRaftCfgMeta meta;
    meta.isStandBy = pSyncInfo->isStandBy;
    meta.snapshotEnable = pSyncInfo->snapshotEnable;
628
    meta.lastConfigIndex = SYNC_INDEX_INVALID;
M
Minghao Li 已提交
629
    ret = raftCfgCreateFile((SSyncCfg*)&(pSyncInfo->syncCfg), meta, pSyncNode->configPath);
M
Minghao Li 已提交
630
    assert(ret == 0);
631 632 633 634 635 636 637

  } else {
    // update syncCfg by raft_config.json
    pSyncNode->pRaftCfg = raftCfgOpen(pSyncNode->configPath);
    assert(pSyncNode->pRaftCfg != NULL);
    pSyncInfo->syncCfg = pSyncNode->pRaftCfg->cfg;

638 639 640 641 642
    if (gRaftDetailLog) {
      char* seralized = raftCfg2Str(pSyncNode->pRaftCfg);
      sInfo("syncNodeOpen update config :%s", seralized);
      taosMemoryFree(seralized);
    }
643 644

    raftCfgClose(pSyncNode->pRaftCfg);
M
Minghao Li 已提交
645 646
  }

M
Minghao Li 已提交
647
  // init by SSyncInfo
M
Minghao Li 已提交
648 649
  pSyncNode->vgId = pSyncInfo->vgId;
  memcpy(pSyncNode->path, pSyncInfo->path, sizeof(pSyncNode->path));
M
Minghao Li 已提交
650
  snprintf(pSyncNode->raftStorePath, sizeof(pSyncNode->raftStorePath), "%s/raft_store.json", pSyncInfo->path);
M
Minghao Li 已提交
651 652
  snprintf(pSyncNode->configPath, sizeof(pSyncNode->configPath), "%s/raft_config.json", pSyncInfo->path);

M
Minghao Li 已提交
653
  pSyncNode->pWal = pSyncInfo->pWal;
S
Shengliang Guan 已提交
654
  pSyncNode->msgcb = pSyncInfo->msgcb;
M
Minghao Li 已提交
655
  pSyncNode->FpSendMsg = pSyncInfo->FpSendMsg;
M
Minghao Li 已提交
656
  pSyncNode->FpEqMsg = pSyncInfo->FpEqMsg;
M
Minghao Li 已提交
657

M
Minghao Li 已提交
658 659 660 661
  // init raft config
  pSyncNode->pRaftCfg = raftCfgOpen(pSyncNode->configPath);
  assert(pSyncNode->pRaftCfg != NULL);

M
Minghao Li 已提交
662
  // init internal
M
Minghao Li 已提交
663 664
  pSyncNode->myNodeInfo = pSyncNode->pRaftCfg->cfg.nodeInfo[pSyncNode->pRaftCfg->cfg.myIndex];
  syncUtilnodeInfo2raftId(&pSyncNode->myNodeInfo, pSyncNode->vgId, &pSyncNode->myRaftId);
M
Minghao Li 已提交
665

M
Minghao Li 已提交
666
  // init peersNum, peers, peersId
M
Minghao Li 已提交
667
  pSyncNode->peersNum = pSyncNode->pRaftCfg->cfg.replicaNum - 1;
M
Minghao Li 已提交
668
  int j = 0;
M
Minghao Li 已提交
669 670 671
  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 已提交
672 673 674
      j++;
    }
  }
M
Minghao Li 已提交
675
  for (int i = 0; i < pSyncNode->peersNum; ++i) {
M
Minghao Li 已提交
676
    syncUtilnodeInfo2raftId(&pSyncNode->peersNodeInfo[i], pSyncNode->vgId, &pSyncNode->peersId[i]);
M
Minghao Li 已提交
677
  }
M
Minghao Li 已提交
678

M
Minghao Li 已提交
679
  // init replicaNum, replicasId
M
Minghao Li 已提交
680 681 682
  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]);
M
Minghao Li 已提交
683 684
  }

M
Minghao Li 已提交
685
  // init raft algorithm
M
Minghao Li 已提交
686
  pSyncNode->pFsm = pSyncInfo->pFsm;
M
Minghao Li 已提交
687
  pSyncNode->quorum = syncUtilQuorum(pSyncNode->pRaftCfg->cfg.replicaNum);
M
Minghao Li 已提交
688 689
  pSyncNode->leaderCache = EMPTY_RAFT_ID;

M
Minghao Li 已提交
690
  // init life cycle outside
M
Minghao Li 已提交
691

M
Minghao Li 已提交
692 693 694 695 696 697 698 699 700 701 702 703 704 705 706 707 708 709 710 711 712 713 714 715
  // 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 已提交
716
  // init TLA+ server vars
M
syncInt  
Minghao Li 已提交
717
  pSyncNode->state = TAOS_SYNC_STATE_FOLLOWER;
M
Minghao Li 已提交
718
  pSyncNode->pRaftStore = raftStoreOpen(pSyncNode->raftStorePath);
M
Minghao Li 已提交
719 720
  assert(pSyncNode->pRaftStore != NULL);

M
Minghao Li 已提交
721
  // init TLA+ candidate vars
M
Minghao Li 已提交
722 723 724 725 726
  pSyncNode->pVotesGranted = voteGrantedCreate(pSyncNode);
  assert(pSyncNode->pVotesGranted != NULL);
  pSyncNode->pVotesRespond = votesRespondCreate(pSyncNode);
  assert(pSyncNode->pVotesRespond != NULL);

M
Minghao Li 已提交
727 728 729 730 731 732 733 734 735
  // init TLA+ leader vars
  pSyncNode->pNextIndex = syncIndexMgrCreate(pSyncNode);
  assert(pSyncNode->pNextIndex != NULL);
  pSyncNode->pMatchIndex = syncIndexMgrCreate(pSyncNode);
  assert(pSyncNode->pMatchIndex != NULL);

  // init TLA+ log vars
  pSyncNode->pLogStore = logStoreCreate(pSyncNode);
  assert(pSyncNode->pLogStore != NULL);
M
Minghao Li 已提交
736
  pSyncNode->commitIndex = SYNC_INDEX_INVALID;
M
Minghao Li 已提交
737

M
Minghao Li 已提交
738 739 740 741 742
  // timer ms init
  pSyncNode->pingBaseLine = PING_TIMER_MS;
  pSyncNode->electBaseLine = ELECT_TIMER_MS_MIN;
  pSyncNode->hbBaseLine = HEARTBEAT_TIMER_MS;

M
Minghao Li 已提交
743
  // init ping timer
M
Minghao Li 已提交
744
  pSyncNode->pPingTimer = NULL;
M
Minghao Li 已提交
745
  pSyncNode->pingTimerMS = pSyncNode->pingBaseLine;
M
Minghao Li 已提交
746 747
  atomic_store_64(&pSyncNode->pingTimerLogicClock, 0);
  atomic_store_64(&pSyncNode->pingTimerLogicClockUser, 0);
M
Minghao Li 已提交
748
  pSyncNode->FpPingTimerCB = syncNodeEqPingTimer;
M
Minghao Li 已提交
749
  pSyncNode->pingTimerCounter = 0;
M
Minghao Li 已提交
750

M
Minghao Li 已提交
751 752
  // init elect timer
  pSyncNode->pElectTimer = NULL;
M
Minghao Li 已提交
753
  pSyncNode->electTimerMS = syncUtilElectRandomMS(pSyncNode->electBaseLine, 2 * pSyncNode->electBaseLine);
M
Minghao Li 已提交
754 755
  atomic_store_64(&pSyncNode->electTimerLogicClock, 0);
  atomic_store_64(&pSyncNode->electTimerLogicClockUser, 0);
M
Minghao Li 已提交
756
  pSyncNode->FpElectTimerCB = syncNodeEqElectTimer;
M
Minghao Li 已提交
757 758 759 760
  pSyncNode->electTimerCounter = 0;

  // init heartbeat timer
  pSyncNode->pHeartbeatTimer = NULL;
M
Minghao Li 已提交
761
  pSyncNode->heartbeatTimerMS = pSyncNode->hbBaseLine;
M
Minghao Li 已提交
762 763
  atomic_store_64(&pSyncNode->heartbeatTimerLogicClock, 0);
  atomic_store_64(&pSyncNode->heartbeatTimerLogicClockUser, 0);
M
Minghao Li 已提交
764
  pSyncNode->FpHeartbeatTimerCB = syncNodeEqHeartbeatTimer;
M
Minghao Li 已提交
765 766
  pSyncNode->heartbeatTimerCounter = 0;

M
Minghao Li 已提交
767
  // init callback
M
Minghao Li 已提交
768 769
  pSyncNode->FpOnPing = syncNodeOnPingCb;
  pSyncNode->FpOnPingReply = syncNodeOnPingReplyCb;
M
Minghao Li 已提交
770
  pSyncNode->FpOnClientRequest = syncNodeOnClientRequestCb;
M
Minghao Li 已提交
771
  pSyncNode->FpOnTimeout = syncNodeOnTimeoutCb;
M
Minghao Li 已提交
772

M
Minghao Li 已提交
773 774 775
  pSyncNode->FpOnSnapshotSend = syncNodeOnSnapshotSendCb;
  pSyncNode->FpOnSnapshotRsp = syncNodeOnSnapshotRspCb;

M
Minghao Li 已提交
776
  if (pSyncNode->pRaftCfg->snapshotEnable) {
M
Minghao Li 已提交
777
    sInfo("sync node use snapshot");
M
Minghao Li 已提交
778 779 780 781
    pSyncNode->FpOnRequestVote = syncNodeOnRequestVoteSnapshotCb;
    pSyncNode->FpOnRequestVoteReply = syncNodeOnRequestVoteReplySnapshotCb;
    pSyncNode->FpOnAppendEntries = syncNodeOnAppendEntriesSnapshotCb;
    pSyncNode->FpOnAppendEntriesReply = syncNodeOnAppendEntriesReplySnapshotCb;
M
Minghao Li 已提交
782 783 784 785 786 787 788

  } else {
    sInfo("sync node do not use snapshot");
    pSyncNode->FpOnRequestVote = syncNodeOnRequestVoteCb;
    pSyncNode->FpOnRequestVoteReply = syncNodeOnRequestVoteReplyCb;
    pSyncNode->FpOnAppendEntries = syncNodeOnAppendEntriesCb;
    pSyncNode->FpOnAppendEntriesReply = syncNodeOnAppendEntriesReplyCb;
M
Minghao Li 已提交
789 790
  }

M
Minghao Li 已提交
791
  // tools
M
Minghao Li 已提交
792
  pSyncNode->pSyncRespMgr = syncRespMgrCreate(pSyncNode, 0);
M
Minghao Li 已提交
793 794
  assert(pSyncNode->pSyncRespMgr != NULL);

795 796
  // restore state
  pSyncNode->restoreFinish = false;
797 798 799 800 801 802

  // pSyncNode->pSnapshot = NULL;
  // if (pSyncNode->pFsm->FpGetSnapshot != NULL) {
  //   pSyncNode->pSnapshot = taosMemoryMalloc(sizeof(SSnapshot));
  //   pSyncNode->pFsm->FpGetSnapshot(pSyncNode->pFsm, pSyncNode->pSnapshot);
  // }
803
  // tsem_init(&(pSyncNode->restoreSem), 0, 0);
804

M
Minghao Li 已提交
805 806 807 808 809 810 811 812
  // 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
813
  pSyncNode->pNewNodeReceiver = snapshotReceiverCreate(pSyncNode, EMPTY_RAFT_ID);
M
Minghao Li 已提交
814

M
Minghao Li 已提交
815
  // start in syncNodeStart
M
Minghao Li 已提交
816
  // start raft
M
Minghao Li 已提交
817
  // syncNodeBecomeFollower(pSyncNode);
M
Minghao Li 已提交
818

819
  // snapshot meta
820
  // pSyncNode->sMeta.lastConfigIndex = -1;
821

822 823
  sDebug("vgId:%d sync event currentTerm:%lu sync open", pSyncNode->vgId, pSyncNode->pRaftStore->currentTerm);

M
Minghao Li 已提交
824 825 826
  return pSyncNode;
}

M
Minghao Li 已提交
827 828
void syncNodeStart(SSyncNode* pSyncNode) {
  // start raft
829
  if (pSyncNode->replicaNum == 1) {
M
Minghao Li 已提交
830
    raftStoreNextTerm(pSyncNode->pRaftStore);
831
    syncNodeBecomeLeader(pSyncNode, "one replica start");
M
format  
Minghao Li 已提交
832

833
    // Raft 3.6.2 Committing entries from previous terms
M
format  
Minghao Li 已提交
834

835 836 837
    // use this now
    syncNodeAppendNoop(pSyncNode);
    syncMaybeAdvanceCommitIndex(pSyncNode);  // maybe only one replica
838

839 840
    if (gRaftDetailLog) {
      syncNodeLog2("==state change become leader immediately==", pSyncNode);
841 842
    }

843 844 845
    return;
  }

846
  syncNodeBecomeFollower(pSyncNode, "first start");
M
Minghao Li 已提交
847

848
  // int32_t ret = 0;
M
Minghao Li 已提交
849
  // ret = syncNodeStartPingTimer(pSyncNode);
850
  // assert(ret == 0);
851

852 853
  if (gRaftDetailLog) {
    syncNodeLog2("==state change become leader immediately==", pSyncNode);
854
  }
M
Minghao Li 已提交
855 856
}

M
Minghao Li 已提交
857 858 859 860 861 862 863 864 865 866 867
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);
}

M
Minghao Li 已提交
868
void syncNodeClose(SSyncNode* pSyncNode) {
869
  sDebug("vgId:%d sync event currentTerm:%lu sync close", pSyncNode->vgId, pSyncNode->pRaftStore->currentTerm);
870

M
Minghao Li 已提交
871
  int32_t ret;
M
Minghao Li 已提交
872
  assert(pSyncNode != NULL);
M
Minghao Li 已提交
873 874 875 876

  ret = raftStoreClose(pSyncNode->pRaftStore);
  assert(ret == 0);

M
Minghao Li 已提交
877
  syncRespMgrDestroy(pSyncNode->pSyncRespMgr);
M
Minghao Li 已提交
878 879 880 881 882
  voteGrantedDestroy(pSyncNode->pVotesGranted);
  votesRespondDestory(pSyncNode->pVotesRespond);
  syncIndexMgrDestroy(pSyncNode->pNextIndex);
  syncIndexMgrDestroy(pSyncNode->pMatchIndex);
  logStoreDestory(pSyncNode->pLogStore);
M
Minghao Li 已提交
883
  raftCfgClose(pSyncNode->pRaftCfg);
M
Minghao Li 已提交
884 885 886 887 888

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

M
Minghao Li 已提交
889 890 891 892
  if (pSyncNode->pFsm != NULL) {
    taosMemoryFree(pSyncNode->pFsm);
  }

M
Minghao Li 已提交
893 894 895 896 897 898 899
  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 已提交
900 901 902 903 904
  if (pSyncNode->pNewNodeReceiver != NULL) {
    snapshotReceiverDestroy(pSyncNode->pNewNodeReceiver);
    pSyncNode->pNewNodeReceiver = NULL;
  }

905
  /*
906 907 908
  if (pSyncNode->pSnapshot != NULL) {
    taosMemoryFree(pSyncNode->pSnapshot);
  }
909
  */
910

911
  // tsem_destroy(&pSyncNode->restoreSem);
912

M
Minghao Li 已提交
913 914
  // free memory in syncFreeNode
  // taosMemoryFree(pSyncNode);
M
Minghao Li 已提交
915 916
}

M
Minghao Li 已提交
917 918 919
// option
bool syncNodeSnapshotEnable(SSyncNode* pSyncNode) { return pSyncNode->pRaftCfg->snapshotEnable; }

M
Minghao Li 已提交
920 921 922 923 924 925 926 927 928 929 930 931 932 933 934
// 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 已提交
935
  SyncPing* pMsg = syncPingBuild3(&pSyncNode->myRaftId, &pSyncNode->myRaftId, pSyncNode->vgId);
M
Minghao Li 已提交
936 937 938 939 940 941 942 943 944 945
  ret = syncNodePing(pSyncNode, &pMsg->destId, pMsg);
  assert(ret == 0);

  syncPingDestroy(pMsg);
  return ret;
}

int32_t syncNodePingPeers(SSyncNode* pSyncNode) {
  int32_t ret = 0;
  for (int i = 0; i < pSyncNode->peersNum; ++i) {
M
Minghao Li 已提交
946 947 948
    SRaftId*  destId = &(pSyncNode->peersId[i]);
    SyncPing* pMsg = syncPingBuild3(&pSyncNode->myRaftId, destId, pSyncNode->vgId);
    ret = syncNodePing(pSyncNode, destId, pMsg);
M
Minghao Li 已提交
949 950 951 952 953 954 955 956
    assert(ret == 0);
    syncPingDestroy(pMsg);
  }
  return ret;
}

int32_t syncNodePingAll(SSyncNode* pSyncNode) {
  int32_t ret = 0;
M
Minghao Li 已提交
957 958 959 960
  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 已提交
961 962 963 964 965 966 967 968 969
    assert(ret == 0);
    syncPingDestroy(pMsg);
  }
  return ret;
}

// timer control --------------
int32_t syncNodeStartPingTimer(SSyncNode* pSyncNode) {
  int32_t ret = 0;
970 971 972 973 974 975 976
  if (syncEnvIsStart()) {
    taosTmrReset(pSyncNode->FpPingTimerCB, pSyncNode->pingTimerMS, pSyncNode, gSyncEnv->pTimerManager,
                 &pSyncNode->pPingTimer);
    atomic_store_64(&pSyncNode->pingTimerLogicClock, pSyncNode->pingTimerLogicClockUser);
  } else {
    sError("sync env is stop, syncNodeStartPingTimer");
  }
M
Minghao Li 已提交
977 978 979 980 981 982 983 984 985 986 987 988 989
  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;
990 991 992 993 994 995 996 997
  if (syncEnvIsStart()) {
    pSyncNode->electTimerMS = ms;
    taosTmrReset(pSyncNode->FpElectTimerCB, pSyncNode->electTimerMS, pSyncNode, gSyncEnv->pTimerManager,
                 &pSyncNode->pElectTimer);
    atomic_store_64(&pSyncNode->electTimerLogicClock, pSyncNode->electTimerLogicClockUser);
  } else {
    sError("sync env is stop, syncNodeStartElectTimer");
  }
M
Minghao Li 已提交
998 999 1000 1001 1002 1003 1004 1005 1006 1007 1008 1009 1010 1011 1012 1013 1014 1015
  return ret;
}

int32_t syncNodeStopElectTimer(SSyncNode* pSyncNode) {
  int32_t ret = 0;
  atomic_add_fetch_64(&pSyncNode->electTimerLogicClockUser, 1);
  taosTmrStop(pSyncNode->pElectTimer);
  pSyncNode->pElectTimer = NULL;
  return ret;
}

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

M
Minghao Li 已提交
1016 1017
int32_t syncNodeResetElectTimer(SSyncNode* pSyncNode) {
  int32_t ret = 0;
M
Minghao Li 已提交
1018 1019 1020 1021 1022 1023 1024
  int32_t electMS;

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

M
Minghao Li 已提交
1029 1030
int32_t syncNodeStartHeartbeatTimer(SSyncNode* pSyncNode) {
  int32_t ret = 0;
1031 1032 1033 1034 1035 1036 1037
  if (syncEnvIsStart()) {
    taosTmrReset(pSyncNode->FpHeartbeatTimerCB, pSyncNode->heartbeatTimerMS, pSyncNode, gSyncEnv->pTimerManager,
                 &pSyncNode->pHeartbeatTimer);
    atomic_store_64(&pSyncNode->heartbeatTimerLogicClock, pSyncNode->heartbeatTimerLogicClockUser);
  } else {
    sError("sync env is stop, syncNodeStartHeartbeatTimer");
  }
M
Minghao Li 已提交
1038 1039 1040 1041 1042 1043 1044 1045 1046 1047 1048 1049 1050 1051 1052
  return ret;
}

int32_t syncNodeStopHeartbeatTimer(SSyncNode* pSyncNode) {
  int32_t ret = 0;
  atomic_add_fetch_64(&pSyncNode->heartbeatTimerLogicClockUser, 1);
  taosTmrStop(pSyncNode->pHeartbeatTimer);
  pSyncNode->pHeartbeatTimer = NULL;
  return ret;
}

// utils --------------
int32_t syncNodeSendMsgById(const SRaftId* destRaftId, SSyncNode* pSyncNode, SRpcMsg* pMsg) {
  SEpSet epSet;
  syncUtilraftId2EpSet(destRaftId, &epSet);
M
Minghao Li 已提交
1053
  if (pSyncNode->FpSendMsg != NULL) {
M
Minghao Li 已提交
1054 1055 1056 1057 1058 1059
    if (gRaftDetailLog) {
      char* JsonStr = syncRpcMsg2Str(pMsg);
      syncUtilJson2Line(JsonStr);
      sTrace("sync send msg, vgId:%d, type:%d, msg:%s", pSyncNode->vgId, pMsg->msgType, JsonStr);
      taosMemoryFree(JsonStr);
    }
M
Minghao Li 已提交
1060

M
Minghao Li 已提交
1061 1062 1063
    // htonl
    syncUtilMsgHtoN(pMsg->pCont);

1064
    pMsg->info.noResp = 1;
S
Shengliang Guan 已提交
1065
    pSyncNode->FpSendMsg(&epSet, pMsg);
M
Minghao Li 已提交
1066 1067 1068
  } else {
    sTrace("syncNodeSendMsgById pSyncNode->FpSendMsg is NULL");
  }
M
Minghao Li 已提交
1069 1070 1071 1072 1073 1074
  return 0;
}

int32_t syncNodeSendMsgByInfo(const SNodeInfo* nodeInfo, SSyncNode* pSyncNode, SRpcMsg* pMsg) {
  SEpSet epSet;
  syncUtilnodeInfo2EpSet(nodeInfo, &epSet);
M
Minghao Li 已提交
1075 1076 1077 1078
  if (pSyncNode->FpSendMsg != NULL) {
    // htonl
    syncUtilMsgHtoN(pMsg->pCont);

1079
    pMsg->info.noResp = 1;
S
Shengliang Guan 已提交
1080
    pSyncNode->FpSendMsg(&epSet, pMsg);
M
Minghao Li 已提交
1081 1082 1083
  } else {
    sTrace("syncNodeSendMsgByInfo pSyncNode->FpSendMsg is NULL");
  }
M
Minghao Li 已提交
1084 1085 1086
  return 0;
}

M
Minghao Li 已提交
1087
cJSON* syncNode2Json(const SSyncNode* pSyncNode) {
C
Cary Xu 已提交
1088
  char   u64buf[128] = {0};
M
Minghao Li 已提交
1089 1090
  cJSON* pRoot = cJSON_CreateObject();

M
Minghao Li 已提交
1091 1092 1093
  if (pSyncNode != NULL) {
    // init by SSyncInfo
    cJSON_AddNumberToObject(pRoot, "vgId", pSyncNode->vgId);
M
Minghao Li 已提交
1094
    cJSON_AddItemToObject(pRoot, "SRaftCfg", raftCfg2Json(pSyncNode->pRaftCfg));
M
Minghao Li 已提交
1095
    cJSON_AddStringToObject(pRoot, "path", pSyncNode->path);
M
Minghao Li 已提交
1096 1097 1098
    cJSON_AddStringToObject(pRoot, "raftStorePath", pSyncNode->raftStorePath);
    cJSON_AddStringToObject(pRoot, "configPath", pSyncNode->configPath);

M
Minghao Li 已提交
1099 1100 1101
    snprintf(u64buf, sizeof(u64buf), "%p", pSyncNode->pWal);
    cJSON_AddStringToObject(pRoot, "pWal", u64buf);

S
Shengliang Guan 已提交
1102
    snprintf(u64buf, sizeof(u64buf), "%p", pSyncNode->msgcb);
M
Minghao Li 已提交
1103 1104 1105 1106
    cJSON_AddStringToObject(pRoot, "rpcClient", u64buf);
    snprintf(u64buf, sizeof(u64buf), "%p", pSyncNode->FpSendMsg);
    cJSON_AddStringToObject(pRoot, "FpSendMsg", u64buf);

S
Shengliang Guan 已提交
1107
    snprintf(u64buf, sizeof(u64buf), "%p", pSyncNode->msgcb);
M
Minghao Li 已提交
1108 1109 1110 1111 1112 1113 1114 1115 1116 1117 1118 1119 1120 1121 1122 1123 1124 1125 1126 1127 1128
    cJSON_AddStringToObject(pRoot, "queue", u64buf);
    snprintf(u64buf, sizeof(u64buf), "%p", pSyncNode->FpEqMsg);
    cJSON_AddStringToObject(pRoot, "FpEqMsg", u64buf);

    // 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 已提交
1129

M
Minghao Li 已提交
1130 1131 1132 1133 1134 1135
    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 已提交
1136

M
Minghao Li 已提交
1137 1138 1139 1140 1141 1142 1143
    // 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 已提交
1144 1145 1146 1147
    // life cycle
    snprintf(u64buf, sizeof(u64buf), "%ld", pSyncNode->rid);
    cJSON_AddStringToObject(pRoot, "rid", u64buf);

M
Minghao Li 已提交
1148 1149 1150
    // tla+ server vars
    cJSON_AddNumberToObject(pRoot, "state", pSyncNode->state);
    cJSON_AddStringToObject(pRoot, "state_str", syncUtilState2String(pSyncNode->state));
M
Minghao Li 已提交
1151
    cJSON_AddItemToObject(pRoot, "pRaftStore", raftStore2Json(pSyncNode->pRaftStore));
M
Minghao Li 已提交
1152 1153 1154 1155 1156 1157 1158 1159 1160 1161 1162

    // 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));
M
Minghao Li 已提交
1163
    snprintf(u64buf, sizeof(u64buf), "%" PRId64 "", pSyncNode->commitIndex);
M
Minghao Li 已提交
1164 1165
    cJSON_AddStringToObject(pRoot, "commitIndex", u64buf);

M
Minghao Li 已提交
1166 1167 1168 1169 1170
    // timer ms init
    cJSON_AddNumberToObject(pRoot, "pingBaseLine", pSyncNode->pingBaseLine);
    cJSON_AddNumberToObject(pRoot, "electBaseLine", pSyncNode->electBaseLine);
    cJSON_AddNumberToObject(pRoot, "hbBaseLine", pSyncNode->hbBaseLine);

M
Minghao Li 已提交
1171 1172 1173 1174
    // ping timer
    snprintf(u64buf, sizeof(u64buf), "%p", pSyncNode->pPingTimer);
    cJSON_AddStringToObject(pRoot, "pPingTimer", u64buf);
    cJSON_AddNumberToObject(pRoot, "pingTimerMS", pSyncNode->pingTimerMS);
M
Minghao Li 已提交
1175
    snprintf(u64buf, sizeof(u64buf), "%" PRIu64 "", pSyncNode->pingTimerLogicClock);
M
Minghao Li 已提交
1176
    cJSON_AddStringToObject(pRoot, "pingTimerLogicClock", u64buf);
M
Minghao Li 已提交
1177
    snprintf(u64buf, sizeof(u64buf), "%" PRIu64 "", pSyncNode->pingTimerLogicClockUser);
M
Minghao Li 已提交
1178 1179 1180
    cJSON_AddStringToObject(pRoot, "pingTimerLogicClockUser", u64buf);
    snprintf(u64buf, sizeof(u64buf), "%p", pSyncNode->FpPingTimerCB);
    cJSON_AddStringToObject(pRoot, "FpPingTimerCB", u64buf);
M
Minghao Li 已提交
1181
    snprintf(u64buf, sizeof(u64buf), "%" PRIu64 "", pSyncNode->pingTimerCounter);
M
Minghao Li 已提交
1182 1183 1184 1185 1186 1187
    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);
M
Minghao Li 已提交
1188
    snprintf(u64buf, sizeof(u64buf), "%" PRIu64 "", pSyncNode->electTimerLogicClock);
M
Minghao Li 已提交
1189
    cJSON_AddStringToObject(pRoot, "electTimerLogicClock", u64buf);
M
Minghao Li 已提交
1190
    snprintf(u64buf, sizeof(u64buf), "%" PRIu64 "", pSyncNode->electTimerLogicClockUser);
M
Minghao Li 已提交
1191 1192 1193
    cJSON_AddStringToObject(pRoot, "electTimerLogicClockUser", u64buf);
    snprintf(u64buf, sizeof(u64buf), "%p", pSyncNode->FpElectTimerCB);
    cJSON_AddStringToObject(pRoot, "FpElectTimerCB", u64buf);
M
Minghao Li 已提交
1194
    snprintf(u64buf, sizeof(u64buf), "%" PRIu64 "", pSyncNode->electTimerCounter);
M
Minghao Li 已提交
1195 1196 1197 1198 1199 1200
    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);
M
Minghao Li 已提交
1201
    snprintf(u64buf, sizeof(u64buf), "%" PRIu64 "", pSyncNode->heartbeatTimerLogicClock);
M
Minghao Li 已提交
1202
    cJSON_AddStringToObject(pRoot, "heartbeatTimerLogicClock", u64buf);
M
Minghao Li 已提交
1203
    snprintf(u64buf, sizeof(u64buf), "%" PRIu64 "", pSyncNode->heartbeatTimerLogicClockUser);
M
Minghao Li 已提交
1204 1205 1206
    cJSON_AddStringToObject(pRoot, "heartbeatTimerLogicClockUser", u64buf);
    snprintf(u64buf, sizeof(u64buf), "%p", pSyncNode->FpHeartbeatTimerCB);
    cJSON_AddStringToObject(pRoot, "FpHeartbeatTimerCB", u64buf);
M
Minghao Li 已提交
1207
    snprintf(u64buf, sizeof(u64buf), "%" PRIu64 "", pSyncNode->heartbeatTimerCounter);
M
Minghao Li 已提交
1208 1209 1210 1211 1212 1213 1214 1215 1216 1217 1218 1219 1220 1221 1222 1223 1224
    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 已提交
1225 1226 1227 1228 1229 1230 1231 1232 1233 1234 1235 1236 1237

    // 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();
1238
    cJSON_AddItemToObject(pRoot, "receiver", snapshotReceiver2Json(pSyncNode->pNewNodeReceiver));
M
Minghao Li 已提交
1239 1240 1241 1242 1243 1244 1245
  }

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

M
Minghao Li 已提交
1246 1247 1248 1249 1250 1251 1252
char* syncNode2Str(const SSyncNode* pSyncNode) {
  cJSON* pJson = syncNode2Json(pSyncNode);
  char*  serialized = cJSON_Print(pJson);
  cJSON_Delete(pJson);
  return serialized;
}

M
Minghao Li 已提交
1253 1254 1255 1256
char* syncNode2SimpleStr(const SSyncNode* pSyncNode) {
  int   len = 256;
  char* s = (char*)taosMemoryMalloc(len);
  snprintf(s, len,
M
Minghao Li 已提交
1257
           "syncNode: vgId:%d currentTerm:%lu, commitIndex:%ld, state:%d %s, isStandBy:%d, "
1258
           "electTimerLogicClock:%lu, "
M
Minghao Li 已提交
1259
           "electTimerLogicClockUser:%lu, "
M
Minghao Li 已提交
1260
           "electTimerMS:%d, replicaNum:%d",
M
Minghao Li 已提交
1261
           pSyncNode->vgId, pSyncNode->pRaftStore->currentTerm, pSyncNode->commitIndex, pSyncNode->state,
1262
           syncUtilState2String(pSyncNode->state), pSyncNode->pRaftCfg->isStandBy, pSyncNode->electTimerLogicClock,
M
Minghao Li 已提交
1263
           pSyncNode->electTimerLogicClockUser, pSyncNode->electTimerMS, pSyncNode->replicaNum);
M
Minghao Li 已提交
1264 1265 1266
  return s;
}

1267 1268 1269 1270 1271 1272 1273 1274 1275 1276 1277 1278 1279 1280 1281 1282 1283 1284 1285 1286 1287 1288 1289 1290 1291 1292 1293 1294
bool syncNodeInConfig(SSyncNode* pSyncNode, const SSyncCfg* config) {
  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;
}

void syncNodeUpdateConfig(SSyncNode* pSyncNode, SSyncCfg* pNewConfig, SyncIndex lastConfigChangeIndex, bool* isDrop) {
1295
  SSyncCfg oldConfig = pSyncNode->pRaftCfg->cfg;
1296
  pSyncNode->pRaftCfg->cfg = *pNewConfig;
1297 1298
  pSyncNode->pRaftCfg->lastConfigIndex = lastConfigChangeIndex;

M
Minghao Li 已提交
1299
  int32_t ret = 0;
M
Minghao Li 已提交
1300

1301 1302 1303 1304
  // save snapshot senders
  int32_t oldReplicaNum = pSyncNode->replicaNum;
  SRaftId oldReplicasId[TSDB_MAX_REPLICA];
  memcpy(oldReplicasId, pSyncNode->replicasId, sizeof(oldReplicasId));
1305
  SSyncSnapshotSender* oldSenders[TSDB_MAX_REPLICA];
1306 1307
  for (int i = 0; i < TSDB_MAX_REPLICA; ++i) {
    oldSenders[i] = (pSyncNode->senders)[i];
1308 1309
    sDebug("vgId:%d sync event currentTerm:%lu save senders %d, %p, privateTerm:%lu", pSyncNode->vgId,
           pSyncNode->pRaftStore->currentTerm, i, oldSenders[i], oldSenders[i]->privateTerm);
1310 1311 1312 1313
    if (gRaftDetailLog) {
      ;
    }
  }
1314

M
Minghao Li 已提交
1315 1316 1317 1318 1319 1320 1321 1322 1323 1324 1325 1326 1327 1328 1329 1330 1331 1332 1333 1334 1335 1336
  // 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]);
  }

  // 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]);
  }
M
Minghao Li 已提交
1337 1338 1339

  syncIndexMgrUpdate(pSyncNode->pNextIndex, pSyncNode);
  syncIndexMgrUpdate(pSyncNode->pMatchIndex, pSyncNode);
M
Minghao Li 已提交
1340 1341
  voteGrantedUpdate(pSyncNode->pVotesGranted, pSyncNode);
  votesRespondUpdate(pSyncNode->pVotesRespond, pSyncNode);
M
Minghao Li 已提交
1342

1343 1344
  pSyncNode->quorum = syncUtilQuorum(pSyncNode->pRaftCfg->cfg.replicaNum);

1345 1346 1347
  // reset snapshot senders

  // clear new
1348 1349 1350
  for (int i = 0; i < TSDB_MAX_REPLICA; ++i) {
    (pSyncNode->senders)[i] = NULL;
  }
1351 1352

  // reset new
1353
  for (int i = 0; i < pSyncNode->replicaNum; ++i) {
1354
    // reset sender
1355
    bool reset = false;
1356 1357
    for (int j = 0; j < TSDB_MAX_REPLICA; ++j) {
      if (syncUtilSameId(&(pSyncNode->replicasId)[i], &oldReplicasId[j])) {
1358
        char     host[128];
1359 1360
        uint16_t port;
        syncUtilU642Addr((pSyncNode->replicasId)[i].addr, host, sizeof(host), &port);
1361 1362 1363
        sDebug("vgId:%d sync event currentTerm:%lu reset sender for %lu, newIndex:%d, %s:%d, %p, privateTerm:%lu",
               pSyncNode->vgId, pSyncNode->pRaftStore->currentTerm, (pSyncNode->replicasId)[i].addr, i, host, port,
               oldSenders[j], oldSenders[j]->privateTerm);
1364
        (pSyncNode->senders)[i] = oldSenders[j];
1365
        oldSenders[j] = NULL;
1366
        reset = true;
1367

1368 1369 1370
        // reset replicaIndex
        int32_t oldreplicaIndex = (pSyncNode->senders)[i]->replicaIndex;
        (pSyncNode->senders)[i]->replicaIndex = i;
1371 1372 1373
        sDebug("vgId:%d sync event currentTerm:%lu udpate replicaIndex from %d to %d, %s:%d, %p, reset:%d",
               pSyncNode->vgId, pSyncNode->pRaftStore->currentTerm, oldreplicaIndex, i, host, port,
               (pSyncNode->senders)[i], reset);
1374 1375 1376
      }
    }
  }
1377 1378

  // create new
1379 1380 1381
  for (int i = 0; i < TSDB_MAX_REPLICA; ++i) {
    if ((pSyncNode->senders)[i] == NULL) {
      (pSyncNode->senders)[i] = snapshotSenderCreate(pSyncNode, i);
1382 1383 1384
      sDebug("vgId:%d sync event currentTerm:%lu create new sender %p replicaIndex:%d, privateTerm:%lu",
             pSyncNode->vgId, pSyncNode->pRaftStore->currentTerm, (pSyncNode->senders)[i], i,
             (pSyncNode->senders)[i]->privateTerm);
1385 1386 1387 1388 1389 1390 1391
    }
  }

  // free old
  for (int i = 0; i < TSDB_MAX_REPLICA; ++i) {
    if (oldSenders[i] != NULL) {
      snapshotSenderDestroy(oldSenders[i]);
1392 1393
      sDebug("vgId:%d sync event currentTerm:%lu delete old sender %p replicaIndex:%d", pSyncNode->vgId,
             pSyncNode->pRaftStore->currentTerm, oldSenders[i], i);
1394
      oldSenders[i] = NULL;
1395 1396 1397
    }
  }

1398 1399 1400 1401
  bool IamInOld = syncNodeInConfig(pSyncNode, &oldConfig);
  bool IamInNew = syncNodeInConfig(pSyncNode, pNewConfig);

#if 0 
1402 1403 1404
  for (int i = 0; i < oldConfig.replicaNum; ++i) {
    if (strcmp((oldConfig.nodeInfo)[i].nodeFqdn, pSyncNode->myNodeInfo.nodeFqdn) == 0 &&
        (oldConfig.nodeInfo)[i].nodePort == pSyncNode->myNodeInfo.nodePort) {
1405
      IamInOld = true;
1406 1407 1408 1409
      break;
    }
  }

M
Minghao Li 已提交
1410 1411 1412
  for (int i = 0; i < newConfig->replicaNum; ++i) {
    if (strcmp((newConfig->nodeInfo)[i].nodeFqdn, pSyncNode->myNodeInfo.nodeFqdn) == 0 &&
        (newConfig->nodeInfo)[i].nodePort == pSyncNode->myNodeInfo.nodePort) {
1413
      IamInNew = true;
M
Minghao Li 已提交
1414 1415 1416
      break;
    }
  }
1417
#endif
M
Minghao Li 已提交
1418

1419 1420 1421 1422 1423
  *isDrop = true;
  if (IamInOld && !IamInNew) {
    *isDrop = true;
  } else {
    *isDrop = false;
M
Minghao Li 已提交
1424 1425
  }

1426 1427 1428 1429
  // may be add me to a new raft group
  if (IamInOld && IamInNew && oldConfig.replicaNum == 1) {
  }

1430 1431 1432
  if (IamInNew) {
    pSyncNode->pRaftCfg->isStandBy = 0;  // change isStandBy to normal
  }
M
Minghao Li 已提交
1433
  raftCfgPersist(pSyncNode->pRaftCfg);
1434 1435 1436 1437

  if (gRaftDetailLog) {
    syncNodeLog2("==syncNodeUpdateConfig==", pSyncNode);
  }
M
Minghao Li 已提交
1438 1439
}

M
Minghao Li 已提交
1440 1441 1442 1443 1444 1445 1446 1447 1448 1449 1450
SSyncNode* syncNodeAcquire(int64_t rid) {
  SSyncNode* pNode = taosAcquireRef(tsNodeRefId, rid);
  if (pNode == NULL) {
    sTrace("failed to acquire node from refId:%" PRId64, rid);
  }

  return pNode;
}

void syncNodeRelease(SSyncNode* pNode) { taosReleaseRef(tsNodeRefId, pNode->rid); }

M
Minghao Li 已提交
1451 1452 1453 1454
// raft state change --------------
void syncNodeUpdateTerm(SSyncNode* pSyncNode, SyncTerm term) {
  if (term > pSyncNode->pRaftStore->currentTerm) {
    raftStoreSetTerm(pSyncNode->pRaftStore, term);
1455 1456 1457
    char tmpBuf[64];
    snprintf(tmpBuf, sizeof(tmpBuf), "update term to %lu", term);
    syncNodeBecomeFollower(pSyncNode, tmpBuf);
M
Minghao Li 已提交
1458 1459 1460 1461
    raftStoreClearVote(pSyncNode->pRaftStore);
  }
}

1462
void syncNodeBecomeFollower(SSyncNode* pSyncNode, const char* debugStr) {
1463 1464 1465 1466 1467
  sDebug(
      "vgId:%d sync event currentTerm:%lu become follower, isStandBy:%d, replicaNum:%d, "
      "restoreFinish:%d, %s",
      pSyncNode->vgId, pSyncNode->pRaftStore->currentTerm, pSyncNode->pRaftCfg->isStandBy, pSyncNode->replicaNum,
      pSyncNode->restoreFinish, debugStr);
M
Minghao Li 已提交
1468

M
Minghao Li 已提交
1469
  // maybe clear leader cache
M
Minghao Li 已提交
1470 1471 1472 1473
  if (pSyncNode->state == TAOS_SYNC_STATE_LEADER) {
    pSyncNode->leaderCache = EMPTY_RAFT_ID;
  }

M
Minghao Li 已提交
1474
  // state change
M
Minghao Li 已提交
1475 1476 1477
  pSyncNode->state = TAOS_SYNC_STATE_FOLLOWER;
  syncNodeStopHeartbeatTimer(pSyncNode);

M
Minghao Li 已提交
1478 1479
  // reset elect timer
  syncNodeResetElectTimer(pSyncNode);
M
Minghao Li 已提交
1480 1481 1482 1483 1484 1485 1486 1487 1488 1489 1490 1491 1492 1493 1494 1495 1496 1497 1498 1499
}

// 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>>
//
1500
void syncNodeBecomeLeader(SSyncNode* pSyncNode, const char* debugStr) {
1501 1502 1503
  // reset restoreFinish
  pSyncNode->restoreFinish = false;

1504 1505
  sDebug("vgId:%d sync event currentTerm:%lu become leader, isStandBy:%d, replicaNum:%d, restoreFinish:%d, %s",
         pSyncNode->vgId, pSyncNode->pRaftStore->currentTerm, pSyncNode->pRaftCfg->isStandBy, pSyncNode->replicaNum,
1506
         pSyncNode->restoreFinish, debugStr);
M
Minghao Li 已提交
1507

M
Minghao Li 已提交
1508
  // state change
M
Minghao Li 已提交
1509
  pSyncNode->state = TAOS_SYNC_STATE_LEADER;
M
Minghao Li 已提交
1510 1511

  // set leader cache
M
Minghao Li 已提交
1512 1513 1514
  pSyncNode->leaderCache = pSyncNode->myRaftId;

  for (int i = 0; i < pSyncNode->pNextIndex->replicaNum; ++i) {
M
Minghao Li 已提交
1515 1516
    // maybe overwrite myself, no harm
    // just do it!
1517 1518 1519 1520 1521 1522 1523 1524 1525

    // 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 已提交
1526 1527 1528
  }

  for (int i = 0; i < pSyncNode->pMatchIndex->replicaNum; ++i) {
M
Minghao Li 已提交
1529 1530
    // maybe overwrite myself, no harm
    // just do it!
M
Minghao Li 已提交
1531 1532 1533
    pSyncNode->pMatchIndex->index[i] = SYNC_INDEX_INVALID;
  }

1534 1535
  // update sender private term
  SSyncSnapshotSender* pMySender = syncNodeGetSnapshotSender(pSyncNode, &(pSyncNode->myRaftId));
1536 1537 1538 1539 1540
  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;
      }
1541
    }
1542
    (pMySender->privateTerm) += 100;
1543 1544
  }

M
Minghao Li 已提交
1545
  // stop elect timer
M
Minghao Li 已提交
1546
  syncNodeStopElectTimer(pSyncNode);
M
Minghao Li 已提交
1547 1548

  // start replicate right now!
M
Minghao Li 已提交
1549
  syncNodeReplicate(pSyncNode);
M
Minghao Li 已提交
1550 1551 1552

  // start heartbeat timer
  syncNodeStartHeartbeatTimer(pSyncNode);
M
Minghao Li 已提交
1553 1554 1555 1556 1557
}

void syncNodeCandidate2Leader(SSyncNode* pSyncNode) {
  assert(pSyncNode->state == TAOS_SYNC_STATE_CANDIDATE);
  assert(voteGrantedMajority(pSyncNode->pVotesGranted));
1558
  syncNodeBecomeLeader(pSyncNode, "candidate to leader");
M
Minghao Li 已提交
1559

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

M
Minghao Li 已提交
1562
  // Raft 3.6.2 Committing entries from previous terms
M
Minghao Li 已提交
1563 1564

  // use this now
M
Minghao Li 已提交
1565
  syncNodeAppendNoop(pSyncNode);
1566
  syncMaybeAdvanceCommitIndex(pSyncNode);  // maybe only one replica
M
Minghao Li 已提交
1567 1568

  // do not use this
M
Minghao Li 已提交
1569
  // syncNodeEqNoop(pSyncNode);
M
Minghao Li 已提交
1570 1571 1572 1573 1574
}

void syncNodeFollower2Candidate(SSyncNode* pSyncNode) {
  assert(pSyncNode->state == TAOS_SYNC_STATE_FOLLOWER);
  pSyncNode->state = TAOS_SYNC_STATE_CANDIDATE;
M
Minghao Li 已提交
1575 1576

  syncNodeLog2("==state change syncNodeFollower2Candidate==", pSyncNode);
M
Minghao Li 已提交
1577 1578 1579 1580
}

void syncNodeLeader2Follower(SSyncNode* pSyncNode) {
  assert(pSyncNode->state == TAOS_SYNC_STATE_LEADER);
1581
  syncNodeBecomeFollower(pSyncNode, "leader to follower");
M
Minghao Li 已提交
1582 1583

  syncNodeLog2("==state change syncNodeLeader2Follower==", pSyncNode);
M
Minghao Li 已提交
1584 1585 1586 1587
}

void syncNodeCandidate2Follower(SSyncNode* pSyncNode) {
  assert(pSyncNode->state == TAOS_SYNC_STATE_CANDIDATE);
1588
  syncNodeBecomeFollower(pSyncNode, "candidate to follower");
M
Minghao Li 已提交
1589 1590

  syncNodeLog2("==state change syncNodeCandidate2Follower==", pSyncNode);
M
Minghao Li 已提交
1591 1592 1593
}

// raft vote --------------
M
Minghao Li 已提交
1594 1595 1596

// just called by syncNodeVoteForSelf
// need assert
M
Minghao Li 已提交
1597 1598 1599 1600 1601 1602 1603
void syncNodeVoteForTerm(SSyncNode* pSyncNode, SyncTerm term, SRaftId* pRaftId) {
  assert(term == pSyncNode->pRaftStore->currentTerm);
  assert(!raftStoreHasVoted(pSyncNode->pRaftStore));

  raftStoreVote(pSyncNode->pRaftStore, pRaftId);
}

M
Minghao Li 已提交
1604
// simulate get vote from outside
M
Minghao Li 已提交
1605 1606 1607
void syncNodeVoteForSelf(SSyncNode* pSyncNode) {
  syncNodeVoteForTerm(pSyncNode, pSyncNode->pRaftStore->currentTerm, &(pSyncNode->myRaftId));

M
Minghao Li 已提交
1608
  SyncRequestVoteReply* pMsg = syncRequestVoteReplyBuild(pSyncNode->vgId);
M
Minghao Li 已提交
1609 1610 1611 1612 1613 1614 1615 1616 1617 1618
  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 已提交
1619
// snapshot --------------
M
Minghao Li 已提交
1620 1621 1622 1623 1624 1625 1626 1627 1628 1629 1630 1631
bool syncNodeHasSnapshot(SSyncNode* pSyncNode) {
  bool      ret = false;
  SSnapshot snapshot = {.data = NULL, .lastApplyIndex = -1, .lastApplyTerm = 0};
  if (pSyncNode->pFsm->FpGetSnapshot != NULL) {
    pSyncNode->pFsm->FpGetSnapshot(pSyncNode->pFsm, &snapshot);
    if (snapshot.lastApplyIndex >= SYNC_INDEX_BEGIN) {
      ret = true;
    }
  }
  return ret;
}

1632
bool syncNodeIsIndexInSnapshot(SSyncNode* pSyncNode, SyncIndex index) {
M
Minghao Li 已提交
1633 1634 1635 1636
  ASSERT(syncNodeHasSnapshot(pSyncNode));
  ASSERT(pSyncNode->pFsm->FpGetSnapshot != NULL);
  ASSERT(index >= SYNC_INDEX_BEGIN);

1637 1638
  SSnapshot snapshot;
  pSyncNode->pFsm->FpGetSnapshot(pSyncNode->pFsm, &snapshot);
M
Minghao Li 已提交
1639
  bool b = (index <= snapshot.lastApplyIndex);
1640 1641 1642
  return b;
}

M
Minghao Li 已提交
1643 1644 1645 1646 1647 1648 1649 1650 1651 1652 1653 1654 1655
SyncIndex syncNodeGetLastIndex(SSyncNode* pSyncNode) {
  SSnapshot snapshot = {.data = NULL, .lastApplyIndex = -1, .lastApplyTerm = 0};
  if (pSyncNode->pFsm->FpGetSnapshot != NULL) {
    pSyncNode->pFsm->FpGetSnapshot(pSyncNode->pFsm, &snapshot);
  }
  SyncIndex logLastIndex = pSyncNode->pLogStore->syncLogLastIndex(pSyncNode->pLogStore);

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

SyncTerm syncNodeGetLastTerm(SSyncNode* pSyncNode) {
  SyncTerm lastTerm = 0;
M
Minghao Li 已提交
1656 1657 1658 1659 1660 1661 1662
  if (syncNodeHasSnapshot(pSyncNode)) {
    // has snapshot
    SSnapshot snapshot = {.data = NULL, .lastApplyIndex = -1, .lastApplyTerm = 0};
    if (pSyncNode->pFsm->FpGetSnapshot != NULL) {
      pSyncNode->pFsm->FpGetSnapshot(pSyncNode->pFsm, &snapshot);
    }

M
Minghao Li 已提交
1663 1664 1665
    SyncIndex logLastIndex = pSyncNode->pLogStore->syncLogLastIndex(pSyncNode->pLogStore);
    if (logLastIndex > snapshot.lastApplyIndex) {
      lastTerm = pSyncNode->pLogStore->syncLogLastTerm(pSyncNode->pLogStore);
M
Minghao Li 已提交
1666 1667 1668 1669
    } else {
      lastTerm = snapshot.lastApplyTerm;
    }

M
Minghao Li 已提交
1670
  } else {
M
Minghao Li 已提交
1671 1672
    // no snapshot
    lastTerm = pSyncNode->pLogStore->syncLogLastTerm(pSyncNode->pLogStore);
1673
  }
M
Minghao Li 已提交
1674

M
Minghao Li 已提交
1675 1676 1677 1678 1679 1680 1681
  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);
1682 1683
  return 0;
}
M
Minghao Li 已提交
1684

M
Minghao Li 已提交
1685 1686 1687 1688 1689
SyncIndex syncNodeSyncStartIndex(SSyncNode* pSyncNode) {
  SyncIndex syncStartIndex = syncNodeGetLastIndex(pSyncNode) + 1;
  return syncStartIndex;
}

M
Minghao Li 已提交
1690
SyncIndex syncNodeGetPreIndex(SSyncNode* pSyncNode, SyncIndex index) {
1691
  ASSERT(index >= SYNC_INDEX_BEGIN);
M
Minghao Li 已提交
1692 1693 1694
  SyncIndex syncStartIndex = syncNodeSyncStartIndex(pSyncNode);
  ASSERT(index <= syncStartIndex);

M
Minghao Li 已提交
1695
  SyncIndex preIndex = index - 1;
M
Minghao Li 已提交
1696 1697 1698
  return preIndex;
}

M
Minghao Li 已提交
1699 1700 1701 1702
SyncTerm syncNodeGetPreTerm(SSyncNode* pSyncNode, SyncIndex index) {
  ASSERT(index >= SYNC_INDEX_BEGIN);
  SyncIndex syncStartIndex = syncNodeSyncStartIndex(pSyncNode);
  ASSERT(index <= syncStartIndex);
M
Minghao Li 已提交
1703

M
Minghao Li 已提交
1704 1705
  if (index == SYNC_INDEX_BEGIN) {
    return 0;
M
Minghao Li 已提交
1706
  }
M
Minghao Li 已提交
1707

M
Minghao Li 已提交
1708 1709 1710 1711 1712 1713 1714 1715 1716 1717 1718 1719 1720 1721 1722
  SyncTerm preTerm = 0;
  if (syncNodeHasSnapshot(pSyncNode)) {
    // has snapshot
    SSnapshot snapshot = {.data = NULL, .lastApplyIndex = -1, .lastApplyTerm = 0};
    if (pSyncNode->pFsm->FpGetSnapshot != NULL) {
      pSyncNode->pFsm->FpGetSnapshot(pSyncNode->pFsm, &snapshot);
    }

    if (index > snapshot.lastApplyIndex + 1) {
      // should be log preTerm
      SSyncRaftEntry* pPreEntry = NULL;
      int32_t         code = pSyncNode->pLogStore->syncLogGetEntry(pSyncNode->pLogStore, index - 1, &pPreEntry);
      ASSERT(code == 0);
      ASSERT(pPreEntry != NULL);

M
Minghao Li 已提交
1723 1724
      preTerm = pPreEntry->term;
      taosMemoryFree(pPreEntry);
M
Minghao Li 已提交
1725 1726 1727

    } else if (index == snapshot.lastApplyIndex + 1) {
      preTerm = snapshot.lastApplyTerm;
M
Minghao Li 已提交
1728

M
Minghao Li 已提交
1729
    } else {
1730
      // maybe snapshot change
M
Minghao Li 已提交
1731 1732 1733 1734 1735 1736 1737 1738 1739 1740
      sError("sync get pre term, bad scene. index:%ld", index);
      logStoreLog2("sync get pre term, bad scene", pSyncNode->pLogStore);

      SSyncRaftEntry* pPreEntry = NULL;
      int32_t         code = pSyncNode->pLogStore->syncLogGetEntry(pSyncNode->pLogStore, index - 1, &pPreEntry);
      ASSERT(code == 0);
      ASSERT(pPreEntry != NULL);

      preTerm = pPreEntry->term;
      taosMemoryFree(pPreEntry);
1741
    }
M
Minghao Li 已提交
1742 1743 1744 1745 1746 1747 1748 1749 1750 1751 1752 1753

  } else {
    // no snapshot
    ASSERT(index > SYNC_INDEX_BEGIN);

    SSyncRaftEntry* pPreEntry = NULL;
    int32_t         code = pSyncNode->pLogStore->syncLogGetEntry(pSyncNode->pLogStore, index - 1, &pPreEntry);
    ASSERT(code == 0);
    ASSERT(pPreEntry != NULL);

    preTerm = pPreEntry->term;
    taosMemoryFree(pPreEntry);
1754
  }
M
Minghao Li 已提交
1755

M
Minghao Li 已提交
1756 1757
  return preTerm;
}
1758

M
Minghao Li 已提交
1759 1760 1761
// 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 已提交
1762
  *pPreTerm = syncNodeGetPreTerm(pSyncNode, index);
M
Minghao Li 已提交
1763 1764 1765
  return 0;
}

M
Minghao Li 已提交
1766 1767 1768 1769 1770
// for debug --------------
void syncNodePrint(SSyncNode* pObj) {
  char* serialized = syncNode2Str(pObj);
  printf("syncNodePrint | len:%lu | %s \n", strlen(serialized), serialized);
  fflush(NULL);
wafwerar's avatar
wafwerar 已提交
1771
  taosMemoryFree(serialized);
M
Minghao Li 已提交
1772 1773 1774 1775 1776 1777
}

void syncNodePrint2(char* s, SSyncNode* pObj) {
  char* serialized = syncNode2Str(pObj);
  printf("syncNodePrint2 | len:%lu | %s | %s \n", strlen(serialized), s, serialized);
  fflush(NULL);
wafwerar's avatar
wafwerar 已提交
1778
  taosMemoryFree(serialized);
M
Minghao Li 已提交
1779 1780 1781 1782
}

void syncNodeLog(SSyncNode* pObj) {
  char* serialized = syncNode2Str(pObj);
M
Minghao Li 已提交
1783
  sTraceLong("syncNodeLog | len:%lu | %s", strlen(serialized), serialized);
wafwerar's avatar
wafwerar 已提交
1784
  taosMemoryFree(serialized);
M
Minghao Li 已提交
1785 1786 1787
}

void syncNodeLog2(char* s, SSyncNode* pObj) {
1788 1789 1790 1791 1792
  if (gRaftDetailLog) {
    char* serialized = syncNode2Str(pObj);
    sTraceLong("syncNodeLog2 | len:%lu | %s | %s", strlen(serialized), s, serialized);
    taosMemoryFree(serialized);
  }
M
Minghao Li 已提交
1793 1794
}

M
Minghao Li 已提交
1795
// ------ local funciton ---------
M
Minghao Li 已提交
1796
// enqueue message ----
M
Minghao Li 已提交
1797 1798
static void syncNodeEqPingTimer(void* param, void* tmrId) {
  SSyncNode* pSyncNode = (SSyncNode*)param;
M
Minghao Li 已提交
1799
  if (atomic_load_64(&pSyncNode->pingTimerLogicClockUser) <= atomic_load_64(&pSyncNode->pingTimerLogicClock)) {
M
Minghao Li 已提交
1800
    SyncTimeout* pSyncMsg = syncTimeoutBuild2(SYNC_TIMEOUT_PING, atomic_load_64(&pSyncNode->pingTimerLogicClock),
M
Minghao Li 已提交
1801
                                              pSyncNode->pingTimerMS, pSyncNode->vgId, pSyncNode);
M
Minghao Li 已提交
1802 1803
    SRpcMsg      rpcMsg;
    syncTimeout2RpcMsg(pSyncMsg, &rpcMsg);
M
Minghao Li 已提交
1804
    syncRpcMsgLog2((char*)"==syncNodeEqPingTimer==", &rpcMsg);
M
Minghao Li 已提交
1805
    if (pSyncNode->FpEqMsg != NULL) {
1806 1807 1808 1809 1810 1811 1812
      int32_t code = pSyncNode->FpEqMsg(pSyncNode->msgcb, &rpcMsg);
      if (code != 0) {
        sError("vgId:%d sync enqueue ping msg error, code:%d", pSyncNode->vgId, code);
        rpcFreeCont(rpcMsg.pCont);
        syncTimeoutDestroy(pSyncMsg);
        return;
      }
M
Minghao Li 已提交
1813 1814 1815
    } else {
      sTrace("syncNodeEqPingTimer pSyncNode->FpEqMsg is NULL");
    }
M
Minghao Li 已提交
1816 1817
    syncTimeoutDestroy(pSyncMsg);

1818 1819 1820 1821 1822 1823 1824
    if (syncEnvIsStart()) {
      taosTmrReset(syncNodeEqPingTimer, pSyncNode->pingTimerMS, pSyncNode, gSyncEnv->pTimerManager,
                   &pSyncNode->pPingTimer);
    } else {
      sError("sync env is stop, syncNodeEqPingTimer");
    }

M
Minghao Li 已提交
1825
  } else {
1826
    sTrace("==syncNodeEqPingTimer== pingTimerLogicClock:%" PRIu64 ", pingTimerLogicClockUser:%" PRIu64 "",
M
Minghao Li 已提交
1827
           pSyncNode->pingTimerLogicClock, pSyncNode->pingTimerLogicClockUser);
M
Minghao Li 已提交
1828 1829 1830 1831 1832 1833 1834
  }
}

static void syncNodeEqElectTimer(void* param, void* tmrId) {
  SSyncNode* pSyncNode = (SSyncNode*)param;
  if (atomic_load_64(&pSyncNode->electTimerLogicClockUser) <= atomic_load_64(&pSyncNode->electTimerLogicClock)) {
    SyncTimeout* pSyncMsg = syncTimeoutBuild2(SYNC_TIMEOUT_ELECTION, atomic_load_64(&pSyncNode->electTimerLogicClock),
M
Minghao Li 已提交
1835
                                              pSyncNode->electTimerMS, pSyncNode->vgId, pSyncNode);
M
Minghao Li 已提交
1836
    SRpcMsg      rpcMsg;
M
Minghao Li 已提交
1837
    syncTimeout2RpcMsg(pSyncMsg, &rpcMsg);
M
Minghao Li 已提交
1838
    syncRpcMsgLog2((char*)"==syncNodeEqElectTimer==", &rpcMsg);
M
Minghao Li 已提交
1839
    if (pSyncNode->FpEqMsg != NULL) {
1840 1841 1842 1843 1844 1845 1846
      int32_t code = pSyncNode->FpEqMsg(pSyncNode->msgcb, &rpcMsg);
      if (code != 0) {
        sError("vgId:%d sync enqueue elect msg error, code:%d", pSyncNode->vgId, code);
        rpcFreeCont(rpcMsg.pCont);
        syncTimeoutDestroy(pSyncMsg);
        return;
      }
M
Minghao Li 已提交
1847
    } else {
1848
      sTrace("syncNodeEqElectTimer FpEqMsg is NULL");
M
Minghao Li 已提交
1849
    }
M
Minghao Li 已提交
1850 1851
    syncTimeoutDestroy(pSyncMsg);

M
Minghao Li 已提交
1852
    // reset timer ms
1853
    if (syncEnvIsStart()) {
1854 1855 1856 1857
      pSyncNode->electTimerMS = syncUtilElectRandomMS(pSyncNode->electBaseLine, 2 * pSyncNode->electBaseLine);
      taosTmrReset(syncNodeEqElectTimer, pSyncNode->electTimerMS, pSyncNode, gSyncEnv->pTimerManager,
                   &pSyncNode->pElectTimer);
    } else {
1858
      sError("sync env is stop, syncNodeEqElectTimer");
1859
    }
M
Minghao Li 已提交
1860
  } else {
1861
    sTrace("==syncNodeEqElectTimer== electTimerLogicClock:%" PRIu64 ", electTimerLogicClockUser:%" PRIu64 "",
M
Minghao Li 已提交
1862
           pSyncNode->electTimerLogicClock, pSyncNode->electTimerLogicClockUser);
M
Minghao Li 已提交
1863 1864 1865
  }
}

M
Minghao Li 已提交
1866 1867
static void syncNodeEqHeartbeatTimer(void* param, void* tmrId) {
  SSyncNode* pSyncNode = (SSyncNode*)param;
1868 1869 1870 1871 1872 1873 1874 1875 1876 1877 1878 1879 1880 1881 1882 1883 1884 1885 1886
  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);
      if (pSyncNode->FpEqMsg != NULL) {
        int32_t code = pSyncNode->FpEqMsg(pSyncNode->msgcb, &rpcMsg);
        if (code != 0) {
          sError("vgId:%d sync enqueue timer msg error, code:%d", pSyncNode->vgId, code);
          rpcFreeCont(rpcMsg.pCont);
          syncTimeoutDestroy(pSyncMsg);
          return;
        }
      } else {
        sError("syncNodeEqHeartbeatTimer FpEqMsg is NULL");
1887
      }
1888
      syncTimeoutDestroy(pSyncMsg);
M
Minghao Li 已提交
1889

1890 1891 1892 1893 1894 1895
      if (syncEnvIsStart()) {
        taosTmrReset(syncNodeEqHeartbeatTimer, pSyncNode->heartbeatTimerMS, pSyncNode, gSyncEnv->pTimerManager,
                     &pSyncNode->pHeartbeatTimer);
      } else {
        sError("sync env is stop, syncNodeEqHeartbeatTimer");
      }
1896
    } else {
1897 1898 1899
      sTrace("==syncNodeEqHeartbeatTimer== heartbeatTimerLogicClock:%" PRIu64 ", heartbeatTimerLogicClockUser:%" PRIu64
             "",
             pSyncNode->heartbeatTimerLogicClock, pSyncNode->heartbeatTimerLogicClockUser);
1900
    }
M
Minghao Li 已提交
1901 1902 1903
  }
}

M
Minghao Li 已提交
1904 1905 1906 1907 1908 1909
static int32_t syncNodeEqNoop(SSyncNode* ths) {
  int32_t ret = 0;
  assert(ths->state == TAOS_SYNC_STATE_LEADER);

  SyncIndex       index = ths->pLogStore->getLastIndex(ths->pLogStore) + 1;
  SyncTerm        term = ths->pRaftStore->currentTerm;
M
Minghao Li 已提交
1910
  SSyncRaftEntry* pEntry = syncEntryBuildNoop(term, index, ths->vgId);
M
Minghao Li 已提交
1911 1912 1913 1914 1915 1916 1917 1918
  assert(pEntry != NULL);

  uint32_t           entryLen;
  char*              serialized = syncEntrySerialize(pEntry, &entryLen);
  SyncClientRequest* pSyncMsg = syncClientRequestBuild(entryLen);
  assert(pSyncMsg->dataLen == entryLen);
  memcpy(pSyncMsg->data, serialized, entryLen);

S
Shengliang Guan 已提交
1919
  SRpcMsg rpcMsg = {0};
M
Minghao Li 已提交
1920
  syncClientRequest2RpcMsg(pSyncMsg, &rpcMsg);
M
Minghao Li 已提交
1921
  if (ths->FpEqMsg != NULL) {
S
Shengliang Guan 已提交
1922
    ths->FpEqMsg(ths->msgcb, &rpcMsg);
M
Minghao Li 已提交
1923 1924 1925
  } else {
    sTrace("syncNodeEqNoop pSyncNode->FpEqMsg is NULL");
  }
M
Minghao Li 已提交
1926

wafwerar's avatar
wafwerar 已提交
1927
  taosMemoryFree(serialized);
M
Minghao Li 已提交
1928 1929 1930 1931 1932 1933 1934 1935 1936 1937
  syncClientRequestDestroy(pSyncMsg);

  return ret;
}

static int32_t syncNodeAppendNoop(SSyncNode* ths) {
  int32_t ret = 0;

  SyncIndex       index = ths->pLogStore->getLastIndex(ths->pLogStore) + 1;
  SyncTerm        term = ths->pRaftStore->currentTerm;
M
Minghao Li 已提交
1938
  SSyncRaftEntry* pEntry = syncEntryBuildNoop(term, index, ths->vgId);
M
Minghao Li 已提交
1939 1940 1941
  assert(pEntry != NULL);

  if (ths->state == TAOS_SYNC_STATE_LEADER) {
1942 1943
    // ths->pLogStore->appendEntry(ths->pLogStore, pEntry);
    ths->pLogStore->syncLogAppendEntry(ths->pLogStore, pEntry);
M
Minghao Li 已提交
1944 1945 1946
    syncNodeReplicate(ths);
  }

M
Minghao Li 已提交
1947
  syncEntryDestory(pEntry);
M
Minghao Li 已提交
1948 1949 1950
  return ret;
}

M
Minghao Li 已提交
1951
// on message ----
M
Minghao Li 已提交
1952 1953
int32_t syncNodeOnPingCb(SSyncNode* ths, SyncPing* pMsg) {
  // log state
1954
  char logBuf[1024] = {0};
M
Minghao Li 已提交
1955 1956 1957 1958 1959 1960
  snprintf(logBuf, sizeof(logBuf),
           "==syncNodeOnPingCb== vgId:%d, state: %d, %s, term:%lu electTimerLogicClock:%lu, "
           "electTimerLogicClockUser:%lu, electTimerMS:%d",
           ths->vgId, ths->state, syncUtilState2String(ths->state), ths->pRaftStore->currentTerm,
           ths->electTimerLogicClock, ths->electTimerLogicClockUser, ths->electTimerMS);

M
Minghao Li 已提交
1961
  int32_t ret = 0;
M
Minghao Li 已提交
1962 1963
  syncPingLog2(logBuf, pMsg);
  SyncPingReply* pMsgReply = syncPingReplyBuild3(&ths->myRaftId, &pMsg->srcId, ths->vgId);
M
Minghao Li 已提交
1964 1965
  SRpcMsg        rpcMsg;
  syncPingReply2RpcMsg(pMsgReply, &rpcMsg);
M
Minghao Li 已提交
1966 1967 1968 1969 1970 1971 1972 1973

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

M
Minghao Li 已提交
1974 1975 1976 1977 1978
  syncNodeSendMsgById(&pMsgReply->destId, ths, &rpcMsg);

  return ret;
}

M
Minghao Li 已提交
1979
int32_t syncNodeOnPingReplyCb(SSyncNode* ths, SyncPingReply* pMsg) {
M
Minghao Li 已提交
1980 1981 1982 1983
  int32_t ret = 0;
  syncPingReplyLog2("==syncNodeOnPingReplyCb==", pMsg);
  return ret;
}
M
Minghao Li 已提交
1984

M
Minghao Li 已提交
1985 1986 1987 1988 1989 1990 1991 1992 1993 1994
// 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 已提交
1995
int32_t syncNodeOnClientRequestCb(SSyncNode* ths, SyncClientRequest* pMsg) {
M
Minghao Li 已提交
1996 1997 1998
  int32_t ret = 0;
  syncClientRequestLog2("==syncNodeOnClientRequestCb==", pMsg);

M
Minghao Li 已提交
1999
  SyncIndex       index = ths->pLogStore->syncLogWriteIndex(ths->pLogStore);
M
Minghao Li 已提交
2000 2001 2002 2003
  SyncTerm        term = ths->pRaftStore->currentTerm;
  SSyncRaftEntry* pEntry = syncEntryBuild2((SyncClientRequest*)pMsg, term, index);
  assert(pEntry != NULL);

M
Minghao Li 已提交
2004
  if (ths->state == TAOS_SYNC_STATE_LEADER) {
2005 2006
    // ths->pLogStore->appendEntry(ths->pLogStore, pEntry);
    ths->pLogStore->syncLogAppendEntry(ths->pLogStore, pEntry);
M
Minghao Li 已提交
2007 2008

    // start replicate right now!
M
Minghao Li 已提交
2009
    syncNodeReplicate(ths);
M
Minghao Li 已提交
2010

M
Minghao Li 已提交
2011 2012 2013 2014
    // pre commit
    SRpcMsg rpcMsg;
    syncEntry2OriginalRpc(pEntry, &rpcMsg);

M
Minghao Li 已提交
2015
    if (ths->pFsm != NULL) {
2016
      // if (ths->pFsm->FpPreCommitCb != NULL && pEntry->originalRpcType != TDMT_SYNC_NOOP) {
M
Minghao Li 已提交
2017
      if (ths->pFsm->FpPreCommitCb != NULL && syncUtilUserPreCommit(pEntry->originalRpcType)) {
M
Minghao Li 已提交
2018 2019 2020 2021 2022 2023 2024
        SFsmCbMeta cbMeta;
        cbMeta.index = pEntry->index;
        cbMeta.isWeak = pEntry->isWeak;
        cbMeta.code = 0;
        cbMeta.state = ths->state;
        cbMeta.seqNum = pEntry->seqNum;
        ths->pFsm->FpPreCommitCb(ths->pFsm, &rpcMsg, cbMeta);
M
Minghao Li 已提交
2025
      }
M
Minghao Li 已提交
2026 2027 2028
    }
    rpcFreeCont(rpcMsg.pCont);

M
Minghao Li 已提交
2029 2030 2031
    // only myself, maybe commit
    syncMaybeAdvanceCommitIndex(ths);

M
Minghao Li 已提交
2032
  } else {
M
Minghao Li 已提交
2033 2034 2035 2036
    // pre commit
    SRpcMsg rpcMsg;
    syncEntry2OriginalRpc(pEntry, &rpcMsg);

M
Minghao Li 已提交
2037
    if (ths->pFsm != NULL) {
2038
      // if (ths->pFsm->FpPreCommitCb != NULL && pEntry->originalRpcType != TDMT_SYNC_NOOP) {
M
Minghao Li 已提交
2039
      if (ths->pFsm->FpPreCommitCb != NULL && syncUtilUserPreCommit(pEntry->originalRpcType)) {
M
Minghao Li 已提交
2040 2041 2042 2043 2044 2045 2046
        SFsmCbMeta cbMeta;
        cbMeta.index = pEntry->index;
        cbMeta.isWeak = pEntry->isWeak;
        cbMeta.code = 1;
        cbMeta.state = ths->state;
        cbMeta.seqNum = pEntry->seqNum;
        ths->pFsm->FpPreCommitCb(ths->pFsm, &rpcMsg, cbMeta);
M
Minghao Li 已提交
2047
      }
M
Minghao Li 已提交
2048 2049
    }
    rpcFreeCont(rpcMsg.pCont);
M
Minghao Li 已提交
2050 2051
  }

M
Minghao Li 已提交
2052
  syncEntryDestory(pEntry);
M
Minghao Li 已提交
2053
  return ret;
2054
}
M
Minghao Li 已提交
2055 2056 2057

static void syncFreeNode(void* param) {
  SSyncNode* pNode = param;
M
Minghao Li 已提交
2058 2059
  // inner object already free
  // syncNodePrint2((char*)"==syncFreeNode==", pNode);
M
Minghao Li 已提交
2060

wafwerar's avatar
wafwerar 已提交
2061
  taosMemoryFree(pNode);
M
Minghao Li 已提交
2062
}
S
Shengliang Guan 已提交
2063 2064 2065 2066

const char* syncStr(ESyncState state) {
  switch (state) {
    case TAOS_SYNC_STATE_FOLLOWER:
2067
      return "follower";
S
Shengliang Guan 已提交
2068
    case TAOS_SYNC_STATE_CANDIDATE:
2069
      return "candidate";
S
Shengliang Guan 已提交
2070
    case TAOS_SYNC_STATE_LEADER:
2071
      return "leader";
S
Shengliang Guan 已提交
2072
    default:
2073
      return "error";
S
Shengliang Guan 已提交
2074
  }
M
Minghao Li 已提交
2075
}
2076

2077
static int32_t syncDoLeaderTransfer(SSyncNode* ths, SRpcMsg* pRpcMsg, SSyncRaftEntry* pEntry) {
M
Minghao Li 已提交
2078 2079
  SyncLeaderTransfer* pSyncLeaderTransfer = syncLeaderTransferFromRpcMsg2(pRpcMsg);

2080
  sDebug("vgId:%d sync event currentTerm:%lu begin leader transfer", ths->vgId, ths->pRaftStore->currentTerm);
M
Minghao Li 已提交
2081

2082 2083
  if (strcmp(pSyncLeaderTransfer->newNodeInfo.nodeFqdn, ths->myNodeInfo.nodeFqdn) == 0 &&
      pSyncLeaderTransfer->newNodeInfo.nodePort == ths->myNodeInfo.nodePort) {
2084 2085 2086
    sDebug("vgId:%d sync event currentTerm:%lu maybe leader transfer to %s:%d %lu", ths->vgId,
           ths->pRaftStore->currentTerm, pSyncLeaderTransfer->newNodeInfo.nodeFqdn,
           pSyncLeaderTransfer->newNodeInfo.nodePort, pSyncLeaderTransfer->newLeaderId.addr);
M
Minghao Li 已提交
2087 2088 2089 2090 2091

    // reset elect timer now!
    int32_t electMS = 1;
    int32_t ret = syncNodeRestartElectTimer(ths, electMS);
    ASSERT(ret == 0);
2092 2093
  }

2094 2095 2096 2097 2098 2099 2100 2101
  /*
    if (syncUtilSameId(&(pSyncLeaderTransfer->newLeaderId), &(ths->myRaftId))) {
      // reset elect timer now!
      int32_t electMS = 1;
      int32_t ret = syncNodeRestartElectTimer(ths, electMS);
      ASSERT(ret == 0);
    }
  */
M
Minghao Li 已提交
2102 2103 2104 2105 2106 2107 2108 2109 2110 2111 2112
  if (ths->pFsm->FpLeaderTransferCb != NULL) {
    SFsmCbMeta cbMeta;
    cbMeta.code = 0;
    cbMeta.currentTerm = ths->pRaftStore->currentTerm;
    cbMeta.flag = 0;
    cbMeta.index = 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);
2113 2114
  }

M
Minghao Li 已提交
2115
  syncLeaderTransferDestroy(pSyncLeaderTransfer);
2116 2117 2118
  return 0;
}

2119 2120 2121 2122 2123 2124 2125 2126 2127 2128 2129 2130 2131 2132 2133
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;
}

2134 2135 2136 2137 2138 2139 2140 2141
static int32_t syncNodeConfigChange(SSyncNode* ths, SRpcMsg* pRpcMsg, SSyncRaftEntry* pEntry) {
  SSyncCfg oldSyncCfg = ths->pRaftCfg->cfg;

  SSyncCfg newSyncCfg;
  int32_t  ret = syncCfgFromStr(pRpcMsg->pCont, &newSyncCfg);
  ASSERT(ret == 0);

  // update new config myIndex
2142 2143 2144 2145 2146 2147 2148 2149 2150 2151 2152 2153 2154 2155
  syncNodeUpdateNewConfigIndex(ths, &newSyncCfg);

  bool IamInNew = syncNodeInConfig(ths, &newSyncCfg);

  /*
   for (int i = 0; i < newSyncCfg.replicaNum; ++i) {
     if (strcmp(ths->myNodeInfo.nodeFqdn, (newSyncCfg.nodeInfo)[i].nodeFqdn) == 0 &&
         ths->myNodeInfo.nodePort == (newSyncCfg.nodeInfo)[i].nodePort) {
       newSyncCfg.myIndex = i;
       IamInNew = true;
       break;
     }
   }
 */
2156 2157 2158

  bool isDrop;

2159
  if (IamInNew) {
2160 2161 2162 2163
    syncNodeUpdateConfig(ths, &newSyncCfg, pEntry->index, &isDrop);

    // change isStandBy to normal
    if (!isDrop) {
2164 2165
      char tmpbuf[128];
      snprintf(tmpbuf, sizeof(tmpbuf), "config change from %d to %d", oldSyncCfg.replicaNum, newSyncCfg.replicaNum);
2166
      if (ths->state == TAOS_SYNC_STATE_LEADER) {
2167
        syncNodeBecomeLeader(ths, tmpbuf);
2168
      } else {
2169
        syncNodeBecomeFollower(ths, tmpbuf);
2170 2171
      }
    }
2172
  } else {
2173 2174 2175
    char tmpbuf[128];
    snprintf(tmpbuf, sizeof(tmpbuf), "config change2 from %d to %d", oldSyncCfg.replicaNum, newSyncCfg.replicaNum);
    syncNodeBecomeFollower(ths, tmpbuf);
2176
  }
2177

2178 2179 2180
  if (gRaftDetailLog) {
    char* sOld = syncCfg2Str(&oldSyncCfg);
    char* sNew = syncCfg2Str(&newSyncCfg);
2181 2182
    sInfo("==config change== 0x11 old:%s new:%s isDrop:%d index:%ld IamInNew:%d \n", sOld, sNew, isDrop, pEntry->index,
          IamInNew);
2183 2184
    taosMemoryFree(sOld);
    taosMemoryFree(sNew);
2185 2186 2187 2188 2189 2190 2191 2192 2193 2194 2195 2196 2197 2198 2199 2200 2201 2202 2203 2204
  }

  // always call FpReConfigCb
  if (ths->pFsm->FpReConfigCb != NULL) {
    SReConfigCbMeta cbMeta = {0};
    cbMeta.code = 0;
    cbMeta.currentTerm = ths->pRaftStore->currentTerm;
    cbMeta.index = pEntry->index;
    cbMeta.term = pEntry->term;
    cbMeta.newCfg = newSyncCfg;
    cbMeta.oldCfg = oldSyncCfg;
    cbMeta.seqNum = pEntry->seqNum;
    cbMeta.flag = 0x11;
    cbMeta.isDrop = isDrop;
    ths->pFsm->FpReConfigCb(ths->pFsm, pRpcMsg, cbMeta);
  }

  return 0;
}

2205
int32_t syncNodeCommit(SSyncNode* ths, SyncIndex beginIndex, SyncIndex endIndex, uint64_t flag) {
2206 2207
  int32_t    code = 0;
  ESyncState state = flag;
2208 2209
  sDebug("vgId:%d sync event currentTerm:%lu commit by wal from index:%" PRId64 " to index:%" PRId64 ", %s", ths->vgId,
         ths->pRaftStore->currentTerm, beginIndex, endIndex, syncUtilState2String(state));
2210 2211 2212 2213 2214 2215 2216 2217 2218 2219 2220 2221 2222

  // execute fsm
  if (ths->pFsm != NULL) {
    for (SyncIndex i = beginIndex; i <= endIndex; ++i) {
      if (i != SYNC_INDEX_INVALID) {
        SSyncRaftEntry* pEntry;
        code = ths->pLogStore->syncLogGetEntry(ths->pLogStore, i, &pEntry);
        ASSERT(code == 0);
        ASSERT(pEntry != NULL);

        SRpcMsg rpcMsg;
        syncEntry2OriginalRpc(pEntry, &rpcMsg);

2223
        // user commit
2224 2225 2226 2227 2228 2229 2230 2231 2232 2233 2234 2235 2236 2237 2238
        if (ths->pFsm->FpCommitCb != NULL && syncUtilUserCommit(pEntry->originalRpcType)) {
          SFsmCbMeta cbMeta;
          cbMeta.index = pEntry->index;
          cbMeta.isWeak = pEntry->isWeak;
          cbMeta.code = 0;
          cbMeta.state = ths->state;
          cbMeta.seqNum = pEntry->seqNum;
          cbMeta.term = pEntry->term;
          cbMeta.currentTerm = ths->pRaftStore->currentTerm;
          cbMeta.flag = flag;

          ths->pFsm->FpCommitCb(ths->pFsm, &rpcMsg, cbMeta);
        }

        // config change
2239
        if (pEntry->originalRpcType == TDMT_SYNC_CONFIG_CHANGE) {
2240 2241 2242
          code = syncNodeConfigChange(ths, &rpcMsg, pEntry);
          ASSERT(code == 0);
        }
2243

M
Minghao Li 已提交
2244
        // leader transfer
2245 2246 2247
        if (pEntry->originalRpcType == TDMT_SYNC_LEADER_TRANSFER) {
          code = syncDoLeaderTransfer(ths, &rpcMsg, pEntry);
          ASSERT(code == 0);
2248 2249 2250 2251 2252 2253 2254 2255 2256
        }

        // restore finish
        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;
2257 2258
            sDebug("vgId:%d sync event currentTerm:%lu restore finish, %s, index:%ld", ths->vgId,
                   ths->pRaftStore->currentTerm, syncUtilState2String(ths->state), pEntry->index);
2259 2260 2261 2262 2263 2264 2265 2266 2267
          }
        }

        rpcFreeCont(rpcMsg.pCont);
        syncEntryDestory(pEntry);
      }
    }
  }
  return 0;
2268 2269 2270 2271 2272 2273 2274 2275 2276
}

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 已提交
2277 2278 2279 2280 2281 2282 2283 2284 2285 2286
}

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 已提交
2287
}