sync.h 7.5 KB
Newer Older
S
Shengliang Guan 已提交
1
/*
M
Minghao Li 已提交
2
 * Copyright (c) 2019 TAOS Data, Inc. <jhtao@taosdata.com>
S
Shengliang Guan 已提交
3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22
 *
 * 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/>.
 */

#ifndef _TD_LIBS_SYNC_H
#define _TD_LIBS_SYNC_H

#ifdef __cplusplus
extern "C" {
#endif

M
Minghao Li 已提交
23
#include "cJSON.h"
M
Minghao Li 已提交
24
#include "tdef.h"
25
#include "tlrucache.h"
S
Shengliang Guan 已提交
26
#include "tmsgcb.h"
S
Shengliang Guan 已提交
27

28 29 30 31 32 33 34
#define SYNC_RESP_TTL_MS             10000000
#define SYNC_SPEED_UP_HB_TIMER       400
#define SYNC_SPEED_UP_AFTER_MS       (1000 * 20)
#define SYNC_SLOW_DOWN_RANGE         100
#define SYNC_MAX_READ_RANGE          2
#define SYNC_MAX_PROGRESS_WAIT_MS    4000
#define SYNC_MAX_START_TIME_RANGE_MS (1000 * 20)
35
#define SYNC_MAX_RECV_TIME_RANGE_MS  1200
M
Minghao Li 已提交
36
#define SYNC_DEL_WAL_MS              (1000 * 60)
37
#define SYNC_ADD_QUORUM_COUNT        3
M
Minghao Li 已提交
38
#define SYNC_MNODE_LOG_RETENTION     10000
39
#define SYNC_VNODE_LOG_RETENTION     100
M
Minghao Li 已提交
40
#define SNAPSHOT_MAX_CLOCK_SKEW_MS   1000 * 10
M
Minghao Li 已提交
41
#define SNAPSHOT_WAIT_MS             1000 * 30
M
Minghao Li 已提交
42

B
Benguang Zhao 已提交
43 44
#define SYNC_MAX_RETRY_BACKOFF         5
#define SYNC_LOG_REPL_RETRY_WAIT_MS    50
M
Minghao Li 已提交
45
#define SYNC_APPEND_ENTRIES_TIMEOUT_MS 10000
M
Minghao Li 已提交
46

M
Minghao Li 已提交
47 48 49
#define SYNC_MAX_BATCH_SIZE 1
#define SYNC_INDEX_BEGIN    0
#define SYNC_INDEX_INVALID  -1
B
Benguang Zhao 已提交
50
#define SYNC_TERM_INVALID   -1  // 0xFFFFFFFFFFFFFFFF
51

M
Minghao Li 已提交
52 53 54 55 56 57
typedef enum {
  SYNC_STRATEGY_NO_SNAPSHOT = 0,
  SYNC_STRATEGY_STANDARD_SNAPSHOT = 1,
  SYNC_STRATEGY_WAL_FIRST = 2,
} ESyncStrategy;

M
Minghao Li 已提交
58
typedef uint64_t SyncNodeId;
S
Shengliang Guan 已提交
59 60
typedef int32_t  SyncGroupId;
typedef int64_t  SyncIndex;
B
Benguang Zhao 已提交
61
typedef int64_t  SyncTerm;
S
Shengliang Guan 已提交
62

63 64 65 66
typedef struct SSyncNode      SSyncNode;
typedef struct SWal           SWal;
typedef struct SSyncRaftEntry SSyncRaftEntry;

S
Shengliang Guan 已提交
67
typedef enum {
M
Minghao Li 已提交
68 69 70
  TAOS_SYNC_STATE_FOLLOWER = 100,
  TAOS_SYNC_STATE_CANDIDATE = 101,
  TAOS_SYNC_STATE_LEADER = 102,
M
Minghao Li 已提交
71
  TAOS_SYNC_STATE_ERROR = 103,
M
syncInt  
Minghao Li 已提交
72
} ESyncState;
S
Shengliang Guan 已提交
73

M
Minghao Li 已提交
74
typedef struct SNodeInfo {
M
Minghao Li 已提交
75 76
  uint16_t nodePort;
  char     nodeFqdn[TSDB_FQDN_LEN];
S
Shengliang Guan 已提交
77 78
} SNodeInfo;

M
Minghao Li 已提交
79
typedef struct SSyncCfg {
M
Minghao Li 已提交
80
  int32_t   replicaNum;
M
Minghao Li 已提交
81
  int32_t   myIndex;
S
Shengliang Guan 已提交
82
  SNodeInfo nodeInfo[TSDB_MAX_REPLICA];
M
Minghao Li 已提交
83
} SSyncCfg;
S
Shengliang Guan 已提交
84

M
Minghao Li 已提交
85
typedef struct SFsmCbMeta {
M
Minghao Li 已提交
86
  int32_t    code;
M
Minghao Li 已提交
87
  SyncIndex  index;
M
Minghao Li 已提交
88 89
  SyncTerm   term;
  uint64_t   seqNum;
90
  SyncIndex  lastConfigIndex;
M
Minghao Li 已提交
91
  ESyncState state;
92
  SyncTerm   currentTerm;
M
Minghao Li 已提交
93
  bool       isWeak;
94
  uint64_t   flag;
M
Minghao Li 已提交
95 96
} SFsmCbMeta;

97
typedef struct SReConfigCbMeta {
M
Minghao Li 已提交
98 99 100 101 102 103 104 105 106 107 108
  int32_t    code;
  SyncIndex  index;
  SyncTerm   term;
  uint64_t   seqNum;
  SyncIndex  lastConfigIndex;
  ESyncState state;
  SyncTerm   currentTerm;
  bool       isWeak;
  uint64_t   flag;

  // config info
M
Minghao Li 已提交
109
  SSyncCfg  oldCfg;
110
  SSyncCfg  newCfg;
M
Minghao Li 已提交
111 112 113 114
  SyncIndex newCfgIndex;
  SyncTerm  newCfgTerm;
  uint64_t  newCfgSeqNum;

115 116
} SReConfigCbMeta;

117 118 119 120 121
typedef struct SSnapshotParam {
  SyncIndex start;
  SyncIndex end;
} SSnapshotParam;

122
typedef struct SSnapshot {
M
Minghao Li 已提交
123
  void*     data;
124 125
  SyncIndex lastApplyIndex;
  SyncTerm  lastApplyTerm;
M
Minghao Li 已提交
126
  SyncIndex lastConfigIndex;
127 128
} SSnapshot;

129 130 131 132
typedef struct SSnapshotMeta {
  SyncIndex lastConfigIndex;
} SSnapshotMeta;

M
Minghao Li 已提交
133 134
typedef struct SSyncFSM {
  void* data;
135

136
  void (*FpCommitCb)(const struct SSyncFSM* pFsm, const SRpcMsg* pMsg, const SFsmCbMeta* pMeta);
S
Shengliang Guan 已提交
137 138
  void (*FpPreCommitCb)(const struct SSyncFSM* pFsm, const SRpcMsg* pMsg, const SFsmCbMeta* pMeta);
  void (*FpRollBackCb)(const struct SSyncFSM* pFsm, const SRpcMsg* pMsg, const SFsmCbMeta* pMeta);
139

S
Shengliang Guan 已提交
140 141 142
  void (*FpRestoreFinishCb)(const struct SSyncFSM* pFsm);
  void (*FpReConfigCb)(const struct SSyncFSM* pFsm, const SRpcMsg* pMsg, const SReConfigCbMeta* pMeta);
  void (*FpLeaderTransferCb)(const struct SSyncFSM* pFsm, const SRpcMsg* pMsg, const SFsmCbMeta* pMeta);
143
  bool (*FpApplyQueueEmptyCb)(const struct SSyncFSM* pFsm);
144
  int32_t (*FpApplyQueueItems)(const struct SSyncFSM* pFsm);
M
Minghao Li 已提交
145

S
Shengliang Guan 已提交
146 147
  void (*FpBecomeLeaderCb)(const struct SSyncFSM* pFsm);
  void (*FpBecomeFollowerCb)(const struct SSyncFSM* pFsm);
148

S
Shengliang Guan 已提交
149 150
  int32_t (*FpGetSnapshot)(const struct SSyncFSM* pFsm, SSnapshot* pSnapshot, void* pReaderParam, void** ppReader);
  int32_t (*FpGetSnapshotInfo)(const struct SSyncFSM* pFsm, SSnapshot* pSnapshot);
151

S
Shengliang Guan 已提交
152 153 154
  int32_t (*FpSnapshotStartRead)(const struct SSyncFSM* pFsm, void* pReaderParam, void** ppReader);
  int32_t (*FpSnapshotStopRead)(const struct SSyncFSM* pFsm, void* pReader);
  int32_t (*FpSnapshotDoRead)(const struct SSyncFSM* pFsm, void* pReader, void** ppBuf, int32_t* len);
155

S
Shengliang Guan 已提交
156 157 158
  int32_t (*FpSnapshotStartWrite)(const struct SSyncFSM* pFsm, void* pWriterParam, void** ppWriter);
  int32_t (*FpSnapshotStopWrite)(const struct SSyncFSM* pFsm, void* pWriter, bool isApply, SSnapshot* pSnapshot);
  int32_t (*FpSnapshotDoWrite)(const struct SSyncFSM* pFsm, void* pWriter, void* pBuf, int32_t len);
159

S
Shengliang Guan 已提交
160 161
} SSyncFSM;

M
Minghao Li 已提交
162 163
// abstract definition of log store in raft
// SWal implements it
S
Shengliang Guan 已提交
164
typedef struct SSyncLogStore {
165 166
  SLRUCache* pCache;
  void*      data;
M
Minghao Li 已提交
167

M
Minghao Li 已提交
168 169
  int32_t (*syncLogUpdateCommitIndex)(struct SSyncLogStore* pLogStore, SyncIndex index);
  SyncIndex (*syncLogCommitIndex)(struct SSyncLogStore* pLogStore);
M
Minghao Li 已提交
170

M
Minghao Li 已提交
171 172
  SyncIndex (*syncLogBeginIndex)(struct SSyncLogStore* pLogStore);
  SyncIndex (*syncLogEndIndex)(struct SSyncLogStore* pLogStore);
M
Minghao Li 已提交
173

M
Minghao Li 已提交
174
  int32_t (*syncLogEntryCount)(struct SSyncLogStore* pLogStore);
175
  int32_t (*syncLogRestoreFromSnapshot)(struct SSyncLogStore* pLogStore, SyncIndex index);
M
Minghao Li 已提交
176
  bool (*syncLogIsEmpty)(struct SSyncLogStore* pLogStore);
177
  bool (*syncLogExist)(struct SSyncLogStore* pLogStore, SyncIndex index);
M
Minghao Li 已提交
178

M
Minghao Li 已提交
179
  SyncIndex (*syncLogWriteIndex)(struct SSyncLogStore* pLogStore);
M
Minghao Li 已提交
180 181 182 183 184 185 186
  SyncIndex (*syncLogLastIndex)(struct SSyncLogStore* pLogStore);
  SyncTerm (*syncLogLastTerm)(struct SSyncLogStore* pLogStore);

  int32_t (*syncLogAppendEntry)(struct SSyncLogStore* pLogStore, SSyncRaftEntry* pEntry);
  int32_t (*syncLogGetEntry)(struct SSyncLogStore* pLogStore, SyncIndex index, SSyncRaftEntry** ppEntry);
  int32_t (*syncLogTruncate)(struct SSyncLogStore* pLogStore, SyncIndex fromIndex);

S
Shengliang Guan 已提交
187 188
} SSyncLogStore;

M
Minghao Li 已提交
189
typedef struct SSyncInfo {
M
Minghao Li 已提交
190 191 192 193 194 195 196 197 198
  bool          isStandBy;
  ESyncStrategy snapshotStrategy;
  SyncGroupId   vgId;
  int32_t       batchSize;
  SSyncCfg      syncCfg;
  char          path[TSDB_FILENAME_LEN];
  SWal*         pWal;
  SSyncFSM*     pFsm;
  SMsgCb*       msgcb;
S
Shengliang Guan 已提交
199 200 201
  int32_t       pingMs;
  int32_t       electMs;
  int32_t       heartbeatMs;
202

S
Shengliang Guan 已提交
203 204 205
  int32_t (*syncSendMSg)(const SEpSet* pEpSet, SRpcMsg* pMsg);
  int32_t (*syncEqMsg)(const SMsgCb* msgcb, SRpcMsg* pMsg);
  int32_t (*syncEqCtrlMsg)(const SMsgCb* msgcb, SRpcMsg* pMsg);
S
Shengliang Guan 已提交
206
} SSyncInfo;
207

208 209 210 211 212 213 214 215
typedef struct SSyncState {
  ESyncState state;
  bool       restored;
} SSyncState;

int32_t syncInit();
void    syncCleanUp();
int64_t syncOpen(SSyncInfo* pSyncInfo);
216
int32_t syncStart(int64_t rid);
217
void    syncStop(int64_t rid);
M
Minghao Li 已提交
218
void    syncPreStop(int64_t rid);
219 220
int32_t syncPropose(int64_t rid, SRpcMsg* pMsg, bool isWeak);
int32_t syncProcessMsg(int64_t rid, SRpcMsg* pMsg);
S
Shengliang Guan 已提交
221
int32_t syncReconfig(int64_t rid, SSyncCfg* pCfg);
222 223
int32_t syncBeginSnapshot(int64_t rid, int64_t lastApplyIndex);
int32_t syncEndSnapshot(int64_t rid);
224
int32_t syncLeaderTransfer(int64_t rid);
M
Minghao Li 已提交
225
int32_t syncStepDown(int64_t rid, SyncTerm newTerm);
226
bool    syncIsReadyForRead(int64_t rid);
M
Minghao Li 已提交
227

228 229 230
SSyncState  syncGetState(int64_t rid);
void        syncGetRetryEpSet(int64_t rid, SEpSet* pEpSet);
const char* syncStr(ESyncState state);
231

S
Shengliang Guan 已提交
232 233 234 235 236
#ifdef __cplusplus
}
#endif

#endif /*_TD_LIBS_SYNC_H*/