syncSnapshot.c 23.8 KB
Newer Older
M
Minghao Li 已提交
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16
/*
 * 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/>.
 */

#include "syncSnapshot.h"
17
#include "syncIndexMgr.h"
18
#include "syncRaftCfg.h"
M
Minghao Li 已提交
19
#include "syncRaftLog.h"
M
Minghao Li 已提交
20 21
#include "syncRaftStore.h"
#include "syncUtil.h"
M
Minghao Li 已提交
22
#include "wal.h"
M
Minghao Li 已提交
23

M
Minghao Li 已提交
24 25
static void snapshotReceiverDoStart(SSyncSnapshotReceiver *pReceiver, SyncTerm privateTerm,
                                    SyncSnapshotSend *pBeginMsg);
M
Minghao Li 已提交
26

27
//----------------------------------
M
Minghao Li 已提交
28
SSyncSnapshotSender *snapshotSenderCreate(SSyncNode *pSyncNode, int32_t replicaIndex) {
M
Minghao Li 已提交
29 30
  bool condition = (pSyncNode->pFsm->FpSnapshotStartRead != NULL) && (pSyncNode->pFsm->FpSnapshotStopRead != NULL) &&
                   (pSyncNode->pFsm->FpSnapshotDoRead != NULL);
M
Minghao Li 已提交
31

M
Minghao Li 已提交
32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47
  SSyncSnapshotSender *pSender = NULL;
  if (condition) {
    pSender = taosMemoryMalloc(sizeof(SSyncSnapshotSender));
    ASSERT(pSender != NULL);
    memset(pSender, 0, sizeof(*pSender));

    pSender->start = false;
    pSender->seq = SYNC_SNAPSHOT_SEQ_INVALID;
    pSender->ack = SYNC_SNAPSHOT_SEQ_INVALID;
    pSender->pReader = NULL;
    pSender->pCurrentBlock = NULL;
    pSender->blockLen = 0;
    pSender->sendingMS = SYNC_SNAPSHOT_RETRY_MS;
    pSender->pSyncNode = pSyncNode;
    pSender->replicaIndex = replicaIndex;
    pSender->term = pSyncNode->pRaftStore->currentTerm;
M
Minghao Li 已提交
48
    pSender->privateTerm = taosGetTimestampMs() + 100;
49
    pSender->pSyncNode->pFsm->FpGetSnapshotInfo(pSender->pSyncNode->pFsm, &(pSender->snapshot));
50
    pSender->finish = false;
M
Minghao Li 已提交
51
  } else {
M
Minghao Li 已提交
52
    sError("snapshotSenderCreate cannot create sender");
M
Minghao Li 已提交
53
  }
54

M
Minghao Li 已提交
55 56
  return pSender;
}
M
Minghao Li 已提交
57

M
Minghao Li 已提交
58 59
void snapshotSenderDestroy(SSyncSnapshotSender *pSender) {
  if (pSender != NULL) {
M
Minghao Li 已提交
60 61 62
    if (pSender->pCurrentBlock != NULL) {
      taosMemoryFree(pSender->pCurrentBlock);
    }
M
Minghao Li 已提交
63 64 65
    taosMemoryFree(pSender);
  }
}
M
Minghao Li 已提交
66

M
Minghao Li 已提交
67 68
bool snapshotSenderIsStart(SSyncSnapshotSender *pSender) { return pSender->start; }

M
Minghao Li 已提交
69
// begin send snapshot (current term, seq begin)
M
Minghao Li 已提交
70
void snapshotSenderStart(SSyncSnapshotSender *pSender, SSnapshot snapshot, void *pReader) {
M
Minghao Li 已提交
71
  ASSERT(!snapshotSenderIsStart(pSender));
M
Minghao Li 已提交
72

M
Minghao Li 已提交
73 74
  pSender->seq = SYNC_SNAPSHOT_SEQ_BEGIN;
  pSender->ack = SYNC_SNAPSHOT_SEQ_INVALID;
M
Minghao Li 已提交
75

76
  // init snapshot and reader
M
Minghao Li 已提交
77
  ASSERT(pSender->pReader == NULL);
M
Minghao Li 已提交
78 79 80
  pSender->pReader = pReader;
  pSender->snapshot = snapshot;

M
Minghao Li 已提交
81 82 83 84 85
  if (pSender->pCurrentBlock != NULL) {
    taosMemoryFree(pSender->pCurrentBlock);
  }
  pSender->blockLen = 0;

86
  if (pSender->snapshot.lastConfigIndex != SYNC_INDEX_INVALID) {
87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115
    int32_t         code = 0;
    SSyncRaftEntry *pEntry = NULL;
    code = pSender->pSyncNode->pLogStore->syncLogGetEntry(pSender->pSyncNode->pLogStore,
                                                          pSender->snapshot.lastConfigIndex, &pEntry);

    bool getLastConfig = false;
    if (code == 0) {
      ASSERT(pEntry != NULL);

      SRpcMsg rpcMsg;
      syncEntry2OriginalRpc(pEntry, &rpcMsg);

      SSyncCfg lastConfig;
      int32_t  ret = syncCfgFromStr(rpcMsg.pCont, &lastConfig);
      ASSERT(ret == 0);
      pSender->lastConfig = lastConfig;
      getLastConfig = true;

      rpcFreeCont(rpcMsg.pCont);
      syncEntryDestory(pEntry);
    } else {
      if (pSender->snapshot.lastConfigIndex == pSender->pSyncNode->pRaftCfg->lastConfigIndex) {
        sTrace("vgId:%d sync sender get cfg from local", pSender->pSyncNode->vgId);
        pSender->lastConfig = pSender->pSyncNode->pRaftCfg->cfg;
        getLastConfig = true;
      }
    }

    if (!getLastConfig) {
116 117 118 119 120 121 122 123 124 125
      char logBuf[128];
      snprintf(logBuf, sizeof(logBuf), "snapshot sender update lcindex from %ld to -1",
               pSender->snapshot.lastConfigIndex);
      pSender->snapshot.lastConfigIndex = -1;

      char *eventLog = snapshotSender2SimpleStr(pSender, logBuf);
      syncNodeEventLog(pSender->pSyncNode, eventLog);
      taosMemoryFree(eventLog);

      memset(&(pSender->lastConfig), 0, sizeof(SSyncCfg));
126
    }
127 128 129 130

  } else {
    memset(&(pSender->lastConfig), 0, sizeof(SSyncCfg));
  }
M
Minghao Li 已提交
131

M
Minghao Li 已提交
132 133 134
  pSender->sendingMS = SYNC_SNAPSHOT_RETRY_MS;
  pSender->term = pSender->pSyncNode->pRaftStore->currentTerm;
  ++(pSender->privateTerm);
135
  pSender->finish = false;
M
Minghao Li 已提交
136 137
  pSender->start = true;

M
Minghao Li 已提交
138
  // build begin msg
M
Minghao Li 已提交
139 140 141 142 143 144
  SyncSnapshotSend *pMsg = syncSnapshotSendBuild(0, pSender->pSyncNode->vgId);
  pMsg->srcId = pSender->pSyncNode->myRaftId;
  pMsg->destId = (pSender->pSyncNode->replicasId)[pSender->replicaIndex];
  pMsg->term = pSender->pSyncNode->pRaftStore->currentTerm;
  pMsg->lastIndex = pSender->snapshot.lastApplyIndex;
  pMsg->lastTerm = pSender->snapshot.lastApplyTerm;
145 146
  pMsg->lastConfigIndex = pSender->snapshot.lastConfigIndex;
  pMsg->lastConfig = pSender->lastConfig;
M
Minghao Li 已提交
147
  pMsg->seq = pSender->seq;  // SYNC_SNAPSHOT_SEQ_BEGIN
M
Minghao Li 已提交
148
  pMsg->privateTerm = pSender->privateTerm;
M
Minghao Li 已提交
149

M
Minghao Li 已提交
150
  // send msg
M
Minghao Li 已提交
151 152 153
  SRpcMsg rpcMsg;
  syncSnapshotSend2RpcMsg(pMsg, &rpcMsg);
  syncNodeSendMsgById(&(pMsg->destId), pSender->pSyncNode, &rpcMsg);
M
Minghao Li 已提交
154

M
Minghao Li 已提交
155 156 157
  char *eventLog = snapshotSender2SimpleStr(pSender, "snapshot sender send");
  syncNodeEventLog(pSender->pSyncNode, eventLog);
  taosMemoryFree(eventLog);
M
Minghao Li 已提交
158

M
Minghao Li 已提交
159 160 161 162
  syncSnapshotSendDestroy(pMsg);
}

void snapshotSenderStop(SSyncSnapshotSender *pSender) {
M
Minghao Li 已提交
163 164 165 166 167
  if (pSender->pReader != NULL) {
    int32_t ret = pSender->pSyncNode->pFsm->FpSnapshotStopRead(pSender->pSyncNode->pFsm, pSender->pReader);
    ASSERT(ret == 0);
    pSender->pReader = NULL;
  }
M
Minghao Li 已提交
168 169 170

  if (pSender->pCurrentBlock != NULL) {
    taosMemoryFree(pSender->pCurrentBlock);
M
Minghao Li 已提交
171
    pSender->pCurrentBlock = NULL;
M
Minghao Li 已提交
172 173
    pSender->blockLen = 0;
  }
M
Minghao Li 已提交
174 175

  pSender->start = false;
M
Minghao Li 已提交
176

177 178 179 180 181
  if (gRaftDetailLog) {
    char *s = snapshotSender2Str(pSender);
    sInfo("snapshotSenderStop %s", s);
    taosMemoryFree(s);
  }
M
Minghao Li 已提交
182 183
}

M
Minghao Li 已提交
184 185
// when sender receiver ack, call this function to send msg from seq
// seq = ack + 1, already updated
M
Minghao Li 已提交
186 187 188 189
int32_t snapshotSend(SSyncSnapshotSender *pSender) {
  // free memory last time (seq - 1)
  if (pSender->pCurrentBlock != NULL) {
    taosMemoryFree(pSender->pCurrentBlock);
M
Minghao Li 已提交
190
    pSender->pCurrentBlock = NULL;
M
Minghao Li 已提交
191 192 193 194 195 196 197
    pSender->blockLen = 0;
  }

  // read data
  int32_t ret = pSender->pSyncNode->pFsm->FpSnapshotDoRead(pSender->pSyncNode->pFsm, pSender->pReader,
                                                           &(pSender->pCurrentBlock), &(pSender->blockLen));
  ASSERT(ret == 0);
M
Minghao Li 已提交
198 199 200 201 202 203
  if (pSender->blockLen > 0) {
    // has read data
  } else {
    // read finish
    pSender->seq = SYNC_SNAPSHOT_SEQ_END;
  }
M
Minghao Li 已提交
204

M
Minghao Li 已提交
205
  // build msg
M
Minghao Li 已提交
206 207 208 209 210 211
  SyncSnapshotSend *pMsg = syncSnapshotSendBuild(pSender->blockLen, pSender->pSyncNode->vgId);
  pMsg->srcId = pSender->pSyncNode->myRaftId;
  pMsg->destId = (pSender->pSyncNode->replicasId)[pSender->replicaIndex];
  pMsg->term = pSender->pSyncNode->pRaftStore->currentTerm;
  pMsg->lastIndex = pSender->snapshot.lastApplyIndex;
  pMsg->lastTerm = pSender->snapshot.lastApplyTerm;
212 213
  pMsg->lastConfigIndex = pSender->snapshot.lastConfigIndex;
  pMsg->lastConfig = pSender->lastConfig;
M
Minghao Li 已提交
214
  pMsg->seq = pSender->seq;
M
Minghao Li 已提交
215
  pMsg->privateTerm = pSender->privateTerm;
M
Minghao Li 已提交
216 217
  memcpy(pMsg->data, pSender->pCurrentBlock, pSender->blockLen);

M
Minghao Li 已提交
218
  // send msg
M
Minghao Li 已提交
219 220 221
  SRpcMsg rpcMsg;
  syncSnapshotSend2RpcMsg(pMsg, &rpcMsg);
  syncNodeSendMsgById(&(pMsg->destId), pSender->pSyncNode, &rpcMsg);
M
Minghao Li 已提交
222 223

  if (pSender->seq == SYNC_SNAPSHOT_SEQ_END) {
M
Minghao Li 已提交
224 225 226 227
    char *eventLog = snapshotSender2SimpleStr(pSender, "snapshot sender finish");
    syncNodeEventLog(pSender->pSyncNode, eventLog);
    taosMemoryFree(eventLog);

M
Minghao Li 已提交
228
  } else {
M
Minghao Li 已提交
229 230 231
    char *eventLog = snapshotSender2SimpleStr(pSender, "snapshot sender sending");
    syncNodeEventLog(pSender->pSyncNode, eventLog);
    taosMemoryFree(eventLog);
M
Minghao Li 已提交
232 233
  }

M
Minghao Li 已提交
234 235 236 237
  syncSnapshotSendDestroy(pMsg);
  return 0;
}

M
Minghao Li 已提交
238 239 240 241 242 243 244 245 246
// send snapshot data from cache
int32_t snapshotReSend(SSyncSnapshotSender *pSender) {
  if (pSender->pCurrentBlock != NULL) {
    SyncSnapshotSend *pMsg = syncSnapshotSendBuild(pSender->blockLen, pSender->pSyncNode->vgId);
    pMsg->srcId = pSender->pSyncNode->myRaftId;
    pMsg->destId = (pSender->pSyncNode->replicasId)[pSender->replicaIndex];
    pMsg->term = pSender->pSyncNode->pRaftStore->currentTerm;
    pMsg->lastIndex = pSender->snapshot.lastApplyIndex;
    pMsg->lastTerm = pSender->snapshot.lastApplyTerm;
247 248
    pMsg->lastConfigIndex = pSender->snapshot.lastConfigIndex;
    pMsg->lastConfig = pSender->lastConfig;
M
Minghao Li 已提交
249 250 251 252 253 254
    pMsg->seq = pSender->seq;
    memcpy(pMsg->data, pSender->pCurrentBlock, pSender->blockLen);

    SRpcMsg rpcMsg;
    syncSnapshotSend2RpcMsg(pMsg, &rpcMsg);
    syncNodeSendMsgById(&(pMsg->destId), pSender->pSyncNode, &rpcMsg);
M
Minghao Li 已提交
255

M
Minghao Li 已提交
256 257 258
    char *eventLog = snapshotSender2SimpleStr(pSender, "snapshot sender resend");
    syncNodeEventLog(pSender->pSyncNode, eventLog);
    taosMemoryFree(eventLog);
M
Minghao Li 已提交
259

M
Minghao Li 已提交
260 261 262 263 264
    syncSnapshotSendDestroy(pMsg);
  }
  return 0;
}

M
Minghao Li 已提交
265 266 267 268 269 270 271 272 273 274 275 276 277 278 279 280 281 282 283 284 285 286 287 288 289 290 291 292
cJSON *snapshotSender2Json(SSyncSnapshotSender *pSender) {
  char   u64buf[128];
  cJSON *pRoot = cJSON_CreateObject();

  if (pSender != NULL) {
    cJSON_AddNumberToObject(pRoot, "start", pSender->start);
    cJSON_AddNumberToObject(pRoot, "seq", pSender->seq);
    cJSON_AddNumberToObject(pRoot, "ack", pSender->ack);

    snprintf(u64buf, sizeof(u64buf), "%p", pSender->pReader);
    cJSON_AddStringToObject(pRoot, "pReader", u64buf);

    snprintf(u64buf, sizeof(u64buf), "%p", pSender->pCurrentBlock);
    cJSON_AddStringToObject(pRoot, "pCurrentBlock", u64buf);
    cJSON_AddNumberToObject(pRoot, "blockLen", pSender->blockLen);

    if (pSender->pCurrentBlock != NULL) {
      char *s;
      s = syncUtilprintBin((char *)(pSender->pCurrentBlock), pSender->blockLen);
      cJSON_AddStringToObject(pRoot, "pCurrentBlock", s);
      taosMemoryFree(s);
      s = syncUtilprintBin2((char *)(pSender->pCurrentBlock), pSender->blockLen);
      cJSON_AddStringToObject(pRoot, "pCurrentBlock2", s);
      taosMemoryFree(s);
    }

    cJSON *pSnapshot = cJSON_CreateObject();
    snprintf(u64buf, sizeof(u64buf), "%lu", pSender->snapshot.lastApplyIndex);
M
Minghao Li 已提交
293
    cJSON_AddStringToObject(pSnapshot, "lastApplyIndex", u64buf);
M
Minghao Li 已提交
294
    snprintf(u64buf, sizeof(u64buf), "%lu", pSender->snapshot.lastApplyTerm);
M
Minghao Li 已提交
295
    cJSON_AddStringToObject(pSnapshot, "lastApplyTerm", u64buf);
M
Minghao Li 已提交
296 297 298 299 300 301 302 303 304
    cJSON_AddItemToObject(pRoot, "snapshot", pSnapshot);

    snprintf(u64buf, sizeof(u64buf), "%lu", pSender->sendingMS);
    cJSON_AddStringToObject(pRoot, "sendingMS", u64buf);
    snprintf(u64buf, sizeof(u64buf), "%p", pSender->pSyncNode);
    cJSON_AddStringToObject(pRoot, "pSyncNode", u64buf);
    cJSON_AddNumberToObject(pRoot, "replicaIndex", pSender->replicaIndex);
    snprintf(u64buf, sizeof(u64buf), "%lu", pSender->term);
    cJSON_AddStringToObject(pRoot, "term", u64buf);
305 306
    snprintf(u64buf, sizeof(u64buf), "%lu", pSender->privateTerm);
    cJSON_AddStringToObject(pRoot, "privateTerm", u64buf);
307
    cJSON_AddNumberToObject(pRoot, "finish", pSender->finish);
M
Minghao Li 已提交
308 309 310 311 312 313 314 315 316
  }

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

char *snapshotSender2Str(SSyncSnapshotSender *pSender) {
  cJSON *pJson = snapshotSender2Json(pSender);
317
  char  *serialized = cJSON_Print(pJson);
M
Minghao Li 已提交
318 319 320 321
  cJSON_Delete(pJson);
  return serialized;
}

M
Minghao Li 已提交
322 323 324 325 326 327 328 329 330 331 332 333 334 335 336 337 338
char *snapshotSender2SimpleStr(SSyncSnapshotSender *pSender, char *event) {
  int32_t len = 256;
  char   *s = taosMemoryMalloc(len);

  SRaftId  destId = pSender->pSyncNode->replicasId[pSender->replicaIndex];
  char     host[128];
  uint16_t port;
  syncUtilU642Addr(destId.addr, host, sizeof(host), &port);

  snprintf(s, len, "%s %p laindex:%ld laterm:%lu lcindex:%ld seq:%d ack:%d finish:%d pterm:%lu replica-index:%d %s:%d",
           event, pSender, pSender->snapshot.lastApplyIndex, pSender->snapshot.lastApplyTerm,
           pSender->snapshot.lastConfigIndex, pSender->seq, pSender->ack, pSender->finish, pSender->privateTerm,
           pSender->replicaIndex, host, port);

  return s;
}

M
Minghao Li 已提交
339
// -------------------------------------
340
SSyncSnapshotReceiver *snapshotReceiverCreate(SSyncNode *pSyncNode, SRaftId fromId) {
M
Minghao Li 已提交
341 342
  bool condition = (pSyncNode->pFsm->FpSnapshotStartWrite != NULL) && (pSyncNode->pFsm->FpSnapshotStopWrite != NULL) &&
                   (pSyncNode->pFsm->FpSnapshotDoWrite != NULL);
343

344
  SSyncSnapshotReceiver *pReceiver = NULL;
M
Minghao Li 已提交
345 346 347 348
  if (condition) {
    pReceiver = taosMemoryMalloc(sizeof(SSyncSnapshotReceiver));
    ASSERT(pReceiver != NULL);
    memset(pReceiver, 0, sizeof(*pReceiver));
349

M
Minghao Li 已提交
350 351 352 353
    pReceiver->start = false;
    pReceiver->ack = SYNC_SNAPSHOT_SEQ_BEGIN;
    pReceiver->pWriter = NULL;
    pReceiver->pSyncNode = pSyncNode;
354
    pReceiver->fromId = fromId;
M
Minghao Li 已提交
355
    pReceiver->term = pSyncNode->pRaftStore->currentTerm;
M
Minghao Li 已提交
356
    pReceiver->privateTerm = 0;
M
Minghao Li 已提交
357 358 359 360
    pReceiver->snapshot.data = NULL;
    pReceiver->snapshot.lastApplyIndex = -1;
    pReceiver->snapshot.lastApplyTerm = 0;
    pReceiver->snapshot.lastConfigIndex = -1;
M
Minghao Li 已提交
361

M
Minghao Li 已提交
362 363 364
  } else {
    sInfo("snapshotReceiverCreate cannot create receiver");
  }
365 366 367

  return pReceiver;
}
M
Minghao Li 已提交
368

369 370 371 372 373
void snapshotReceiverDestroy(SSyncSnapshotReceiver *pReceiver) {
  if (pReceiver != NULL) {
    taosMemoryFree(pReceiver);
  }
}
M
Minghao Li 已提交
374

375 376
bool snapshotReceiverIsStart(SSyncSnapshotReceiver *pReceiver) { return pReceiver->start; }

M
Minghao Li 已提交
377
// begin receive snapshot msg (current term, seq begin)
M
Minghao Li 已提交
378 379
static void snapshotReceiverDoStart(SSyncSnapshotReceiver *pReceiver, SyncTerm privateTerm,
                                    SyncSnapshotSend *pBeginMsg) {
M
Minghao Li 已提交
380
  pReceiver->term = pReceiver->pSyncNode->pRaftStore->currentTerm;
381
  pReceiver->privateTerm = privateTerm;
M
Minghao Li 已提交
382
  pReceiver->ack = SYNC_SNAPSHOT_SEQ_BEGIN;
M
Minghao Li 已提交
383 384 385 386 387
  pReceiver->fromId = pBeginMsg->srcId;

  pReceiver->snapshot.lastApplyIndex = pBeginMsg->lastIndex;
  pReceiver->snapshot.lastApplyTerm = pBeginMsg->lastTerm;
  pReceiver->snapshot.lastConfigIndex = pBeginMsg->lastConfigIndex;
M
Minghao Li 已提交
388 389 390 391 392 393 394 395

  ASSERT(pReceiver->pWriter == NULL);
  int32_t ret = pReceiver->pSyncNode->pFsm->FpSnapshotStartWrite(pReceiver->pSyncNode->pFsm, &(pReceiver->pWriter));
  ASSERT(ret == 0);
}

// if receiver receive msg from seq = SYNC_SNAPSHOT_SEQ_BEGIN, start receiver
// if already start, force close, start again
M
Minghao Li 已提交
396
void snapshotReceiverStart(SSyncSnapshotReceiver *pReceiver, SyncTerm privateTerm, SyncSnapshotSend *pBeginMsg) {
397
  if (!snapshotReceiverIsStart(pReceiver)) {
M
Minghao Li 已提交
398
    // start
M
Minghao Li 已提交
399
    snapshotReceiverDoStart(pReceiver, privateTerm, pBeginMsg);
M
Minghao Li 已提交
400
    pReceiver->start = true;
M
Minghao Li 已提交
401

402
  } else {
M
Minghao Li 已提交
403
    // already start
404
    sInfo("snapshot recv, receiver already start");
M
Minghao Li 已提交
405 406 407 408 409 410 411 412

    // force close, abandon incomplete data
    int32_t ret =
        pReceiver->pSyncNode->pFsm->FpSnapshotStopWrite(pReceiver->pSyncNode->pFsm, pReceiver->pWriter, false);
    ASSERT(ret == 0);
    pReceiver->pWriter = NULL;

    // start again
M
Minghao Li 已提交
413
    snapshotReceiverDoStart(pReceiver, privateTerm, pBeginMsg);
M
Minghao Li 已提交
414
    pReceiver->start = true;
415
  }
M
Minghao Li 已提交
416

417 418 419 420 421
  if (gRaftDetailLog) {
    char *s = snapshotReceiver2Str(pReceiver);
    sInfo("snapshotReceiverStart %s", s);
    taosMemoryFree(s);
  }
422
}
M
Minghao Li 已提交
423

424
void snapshotReceiverStop(SSyncSnapshotReceiver *pReceiver, bool apply) {
M
Minghao Li 已提交
425 426 427 428 429
  if (pReceiver->pWriter != NULL) {
    int32_t ret =
        pReceiver->pSyncNode->pFsm->FpSnapshotStopWrite(pReceiver->pSyncNode->pFsm, pReceiver->pWriter, false);
    ASSERT(ret == 0);
    pReceiver->pWriter = NULL;
430
  }
M
Minghao Li 已提交
431 432

  pReceiver->start = false;
433 434

  if (apply) {
435
    //    ++(pReceiver->privateTerm);
436
  }
M
Minghao Li 已提交
437

438 439 440 441 442
  if (gRaftDetailLog) {
    char *s = snapshotReceiver2Str(pReceiver);
    sInfo("snapshotReceiverStop %s", s);
    taosMemoryFree(s);
  }
443 444 445 446 447
}

cJSON *snapshotReceiver2Json(SSyncSnapshotReceiver *pReceiver) {
  char   u64buf[128];
  cJSON *pRoot = cJSON_CreateObject();
M
Minghao Li 已提交
448

449 450 451
  if (pReceiver != NULL) {
    cJSON_AddNumberToObject(pRoot, "start", pReceiver->start);
    cJSON_AddNumberToObject(pRoot, "ack", pReceiver->ack);
M
Minghao Li 已提交
452

453 454 455 456 457
    snprintf(u64buf, sizeof(u64buf), "%p", pReceiver->pWriter);
    cJSON_AddStringToObject(pRoot, "pWriter", u64buf);

    snprintf(u64buf, sizeof(u64buf), "%p", pReceiver->pSyncNode);
    cJSON_AddStringToObject(pRoot, "pSyncNode", u64buf);
458 459 460 461 462 463

    cJSON *pFromId = cJSON_CreateObject();
    snprintf(u64buf, sizeof(u64buf), "%lu", pReceiver->fromId.addr);
    cJSON_AddStringToObject(pFromId, "addr", u64buf);
    {
      uint64_t u64 = pReceiver->fromId.addr;
464
      cJSON   *pTmp = pFromId;
465 466 467 468 469 470 471 472 473
      char     host[128] = {0};
      uint16_t port;
      syncUtilU642Addr(u64, host, sizeof(host), &port);
      cJSON_AddStringToObject(pTmp, "addr_host", host);
      cJSON_AddNumberToObject(pTmp, "addr_port", port);
    }
    cJSON_AddNumberToObject(pFromId, "vgId", pReceiver->fromId.vgId);
    cJSON_AddItemToObject(pRoot, "fromId", pFromId);

M
Minghao Li 已提交
474 475 476 477 478 479 480 481 482
    snprintf(u64buf, sizeof(u64buf), "%lu", pReceiver->snapshot.lastApplyIndex);
    cJSON_AddStringToObject(pRoot, "snapshot.lastApplyIndex", u64buf);

    snprintf(u64buf, sizeof(u64buf), "%lu", pReceiver->snapshot.lastApplyTerm);
    cJSON_AddStringToObject(pRoot, "snapshot.lastApplyTerm", u64buf);

    snprintf(u64buf, sizeof(u64buf), "%lu", pReceiver->snapshot.lastConfigIndex);
    cJSON_AddStringToObject(pRoot, "snapshot.lastConfigIndex", u64buf);

483 484
    snprintf(u64buf, sizeof(u64buf), "%lu", pReceiver->term);
    cJSON_AddStringToObject(pRoot, "term", u64buf);
485 486 487

    snprintf(u64buf, sizeof(u64buf), "%lu", pReceiver->privateTerm);
    cJSON_AddStringToObject(pRoot, "privateTerm", u64buf);
488 489 490 491 492 493 494 495 496
  }

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

char *snapshotReceiver2Str(SSyncSnapshotReceiver *pReceiver) {
  cJSON *pJson = snapshotReceiver2Json(pReceiver);
497
  char  *serialized = cJSON_Print(pJson);
498 499 500
  cJSON_Delete(pJson);
  return serialized;
}
M
Minghao Li 已提交
501

M
Minghao Li 已提交
502 503 504 505 506 507 508 509 510
char *snapshotReceiver2SimpleStr(SSyncSnapshotReceiver *pReceiver, char *event) {
  int32_t len = 256;
  char   *s = taosMemoryMalloc(len);

  SRaftId  fromId = pReceiver->fromId;
  char     host[128];
  uint16_t port;
  syncUtilU642Addr(fromId.addr, host, sizeof(host), &port);

M
Minghao Li 已提交
511 512 513
  snprintf(s, len, "%s %p start:%d ack:%d term:%lu pterm:%lu from:%s:%d laindex:%ld laterm:%lu lcindex:%ld", event,
           pReceiver, pReceiver->start, pReceiver->ack, pReceiver->term, pReceiver->privateTerm, host, port,
           pReceiver->snapshot.lastApplyIndex, pReceiver->snapshot.lastApplyTerm, pReceiver->snapshot.lastConfigIndex);
M
Minghao Li 已提交
514 515 516 517

  return s;
}

518
// receiver do something
M
Minghao Li 已提交
519
int32_t syncNodeOnSnapshotSendCb(SSyncNode *pSyncNode, SyncSnapshotSend *pMsg) {
M
Minghao Li 已提交
520
  // get receiver
521 522
  SSyncSnapshotReceiver *pReceiver = pSyncNode->pNewNodeReceiver;
  bool                   needRsp = false;
523
  int32_t                writeCode = 0;
M
Minghao Li 已提交
524

525 526
  // state, term, seq/ack
  if (pSyncNode->state == TAOS_SYNC_STATE_FOLLOWER) {
M
Minghao Li 已提交
527 528 529
    if (pMsg->term == pSyncNode->pRaftStore->currentTerm) {
      if (pMsg->seq == SYNC_SNAPSHOT_SEQ_BEGIN) {
        // begin
M
Minghao Li 已提交
530
        snapshotReceiverStart(pReceiver, pMsg->privateTerm, pMsg);
M
Minghao Li 已提交
531 532 533
        pReceiver->ack = pMsg->seq;
        needRsp = true;

M
Minghao Li 已提交
534 535 536
        char *eventLog = snapshotReceiver2SimpleStr(pReceiver, "snapshot receiver begin");
        syncNodeEventLog(pSyncNode, eventLog);
        taosMemoryFree(eventLog);
M
Minghao Li 已提交
537

M
Minghao Li 已提交
538 539
      } else if (pMsg->seq == SYNC_SNAPSHOT_SEQ_END) {
        // end, finish FSM
540 541 542
        writeCode = pSyncNode->pFsm->FpSnapshotDoWrite(pSyncNode->pFsm, pReceiver->pWriter, pMsg->data, pMsg->dataLen);
        ASSERT(writeCode == 0);

M
Minghao Li 已提交
543
        pSyncNode->pFsm->FpSnapshotStopWrite(pSyncNode->pFsm, pReceiver->pWriter, true);
544 545 546 547
        if (pReceiver->snapshot.lastApplyIndex > pReceiver->pSyncNode->commitIndex) {
          pReceiver->pSyncNode->commitIndex = pReceiver->snapshot.lastApplyIndex;
        }

548 549
        // pSyncNode->pLogStore->syncLogSetBeginIndex(pSyncNode->pLogStore, pMsg->lastIndex + 1);
        pSyncNode->pLogStore->syncLogRestoreFromSnapshot(pSyncNode->pLogStore, pMsg->lastIndex);
550

551 552
        // maybe update lastconfig
        if (pMsg->lastConfigIndex >= SYNC_INDEX_BEGIN) {
M
Minghao Li 已提交
553 554
          // int32_t  oldReplicaNum = pSyncNode->replicaNum;
          SSyncCfg oldSyncCfg = pSyncNode->pRaftCfg->cfg;
555

556 557 558
          // update new config myIndex
          SSyncCfg newSyncCfg = pMsg->lastConfig;
          syncNodeUpdateNewConfigIndex(pSyncNode, &newSyncCfg);
M
Minghao Li 已提交
559 560 561

          // do config change
          syncNodeDoConfigChange(pSyncNode, &newSyncCfg, pMsg->lastConfigIndex);
562 563
        }

M
Minghao Li 已提交
564
        SSnapshot snapshot;
565
        pSyncNode->pFsm->FpGetSnapshotInfo(pSyncNode->pFsm, &snapshot);
566

M
Minghao Li 已提交
567
        do {
568
          char *eventLog = snapshotReceiver2SimpleStr(pReceiver, "snapshot receiver finish, apply snapshot");
M
Minghao Li 已提交
569 570 571
          syncNodeEventLog(pSyncNode, eventLog);
          taosMemoryFree(eventLog);
        } while (0);
M
Minghao Li 已提交
572

M
Minghao Li 已提交
573
        pReceiver->pWriter = NULL;
574
        snapshotReceiverStop(pReceiver, true);
M
Minghao Li 已提交
575
        pReceiver->ack = pMsg->seq;
M
Minghao Li 已提交
576
        needRsp = true;
M
Minghao Li 已提交
577

M
Minghao Li 已提交
578
        do {
579
          char *eventLog = snapshotReceiver2SimpleStr(pReceiver, "snapshot receiver stop");
M
Minghao Li 已提交
580 581 582
          syncNodeEventLog(pSyncNode, eventLog);
          taosMemoryFree(eventLog);
        } while (0);
M
Minghao Li 已提交
583

M
Minghao Li 已提交
584 585
      } else if (pMsg->seq == SYNC_SNAPSHOT_SEQ_FORCE_CLOSE) {
        pSyncNode->pFsm->FpSnapshotStopWrite(pSyncNode->pFsm, pReceiver->pWriter, false);
586
        snapshotReceiverStop(pReceiver, false);
M
Minghao Li 已提交
587 588
        needRsp = false;

M
Minghao Li 已提交
589 590 591 592 593
        do {
          char *eventLog = snapshotReceiver2SimpleStr(pReceiver, "snapshot receiver force close");
          syncNodeEventLog(pSyncNode, eventLog);
          taosMemoryFree(eventLog);
        } while (0);
M
Minghao Li 已提交
594

M
Minghao Li 已提交
595 596 597
      } else if (pMsg->seq > SYNC_SNAPSHOT_SEQ_BEGIN && pMsg->seq < SYNC_SNAPSHOT_SEQ_END) {
        // transfering
        if (pMsg->seq == pReceiver->ack + 1) {
598 599 600
          writeCode =
              pSyncNode->pFsm->FpSnapshotDoWrite(pSyncNode->pFsm, pReceiver->pWriter, pMsg->data, pMsg->dataLen);
          ASSERT(writeCode == 0);
M
Minghao Li 已提交
601 602 603 604
          pReceiver->ack = pMsg->seq;
        }
        needRsp = true;

M
Minghao Li 已提交
605 606 607 608 609
        do {
          char *eventLog = snapshotReceiver2SimpleStr(pReceiver, "snapshot receiver receiving");
          syncNodeEventLog(pSyncNode, eventLog);
          taosMemoryFree(eventLog);
        } while (0);
M
Minghao Li 已提交
610

M
Minghao Li 已提交
611 612
      } else {
        ASSERT(0);
613
      }
M
Minghao Li 已提交
614

M
Minghao Li 已提交
615 616 617 618 619 620 621 622
      if (needRsp) {
        SyncSnapshotRsp *pRspMsg = syncSnapshotRspBuild(pSyncNode->vgId);
        pRspMsg->srcId = pSyncNode->myRaftId;
        pRspMsg->destId = pMsg->srcId;
        pRspMsg->term = pSyncNode->pRaftStore->currentTerm;
        pRspMsg->lastIndex = pMsg->lastIndex;
        pRspMsg->lastTerm = pMsg->lastTerm;
        pRspMsg->ack = pReceiver->ack;
623
        pRspMsg->code = writeCode;
M
Minghao Li 已提交
624
        pRspMsg->privateTerm = pReceiver->privateTerm;
M
Minghao Li 已提交
625 626 627 628

        SRpcMsg rpcMsg;
        syncSnapshotRsp2RpcMsg(pRspMsg, &rpcMsg);
        syncNodeSendMsgById(&(pRspMsg->destId), pSyncNode, &rpcMsg);
M
Minghao Li 已提交
629

M
Minghao Li 已提交
630 631
        syncSnapshotRspDestroy(pRspMsg);
      }
M
Minghao Li 已提交
632
    }
M
Minghao Li 已提交
633 634
  } else {
    syncNodeLog2("syncNodeOnSnapshotSendCb not follower", pSyncNode);
M
Minghao Li 已提交
635
  }
M
Minghao Li 已提交
636

M
Minghao Li 已提交
637 638 639
  return 0;
}

M
Minghao Li 已提交
640 641
// sender receives ack, set seq = ack + 1, send msg from seq
// if ack == SYNC_SNAPSHOT_SEQ_END, stop sender
642
int32_t syncNodeOnSnapshotRspCb(SSyncNode *pSyncNode, SyncSnapshotRsp *pMsg) {
643 644 645 646 647 648
  // if already drop replica, do not process
  if (!syncNodeInRaftGroup(pSyncNode, &(pMsg->srcId)) && pSyncNode->state == TAOS_SYNC_STATE_LEADER) {
    sInfo("recv SyncSnapshotRsp maybe replica already dropped");
    return 0;
  }

M
Minghao Li 已提交
649
  // get sender
650
  SSyncSnapshotSender *pSender = syncNodeGetSnapshotSender(pSyncNode, &(pMsg->srcId));
651 652 653 654 655
  ASSERT(pSender != NULL);

  // state, term, seq/ack
  if (pSyncNode->state == TAOS_SYNC_STATE_LEADER) {
    if (pMsg->term == pSyncNode->pRaftStore->currentTerm) {
M
Minghao Li 已提交
656 657
      // receiver ack is finish, close sender
      if (pMsg->ack == SYNC_SNAPSHOT_SEQ_END) {
658
        pSender->finish = true;
M
Minghao Li 已提交
659 660 661 662 663
        snapshotSenderStop(pSender);
        return 0;
      }

      // send next msg
664
      if (pMsg->ack == pSender->seq) {
M
Minghao Li 已提交
665
        // update sender ack
666 667
        pSender->ack = pMsg->ack;
        (pSender->seq)++;
M
Minghao Li 已提交
668
        snapshotSend(pSender);
669

M
Minghao Li 已提交
670 671
      } else if (pMsg->ack == pSender->seq - 1) {
        snapshotReSend(pSender);
672

M
Minghao Li 已提交
673 674
      } else {
        ASSERT(0);
675 676
      }
    }
M
Minghao Li 已提交
677 678
  } else {
    syncNodeLog2("syncNodeOnSnapshotRspCb not leader", pSyncNode);
679 680 681
  }

  return 0;
682
}