syncSnapshot.c 23.7 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;
        }

M
Minghao Li 已提交
548
        pSyncNode->pLogStore->syncLogSetBeginIndex(pSyncNode->pLogStore, pMsg->lastIndex + 1);
549

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

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

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

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

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

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

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

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

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

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

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

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

M
Minghao Li 已提交
614 615 616 617 618 619 620 621
      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;
622
        pRspMsg->code = writeCode;
M
Minghao Li 已提交
623
        pRspMsg->privateTerm = pReceiver->privateTerm;
M
Minghao Li 已提交
624 625 626 627

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

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

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

M
Minghao Li 已提交
639 640
// sender receives ack, set seq = ack + 1, send msg from seq
// if ack == SYNC_SNAPSHOT_SEQ_END, stop sender
641
int32_t syncNodeOnSnapshotRspCb(SSyncNode *pSyncNode, SyncSnapshotRsp *pMsg) {
642 643 644 645 646 647
  // 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 已提交
648
  // get sender
649
  SSyncSnapshotSender *pSender = syncNodeGetSnapshotSender(pSyncNode, &(pMsg->srcId));
650 651 652 653 654
  ASSERT(pSender != NULL);

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

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

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

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

  return 0;
681
}