mndDef.h 18.3 KB
Newer Older
H
refact  
Hongze Cheng 已提交
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15
/*
 * Copyright (c) 2019 TAOS Data, Inc. <jhtao@taosdata.com>
 *
 * This program is free software: you can use, redistribute, and/or modify
 * it under the terms of the GNU Affero General Public License, version 3
 * or later ("AGPL"), as published by the Free Software Foundation.
 *
 * This program is distributed in the hope that it will be useful, but WITHOUT
 * ANY WARRANTY; without even the implied warranty of MERCHANTABILITY or
 * FITNESS FOR A PARTICULAR PURPOSE.
 *
 * You should have received a copy of the GNU Affero General Public License
 * along with this program. If not, see <http://www.gnu.org/licenses/>.
 */

S
Shengliang Guan 已提交
16 17
#ifndef _TD_MND_DEF_H_
#define _TD_MND_DEF_H_
S
Shengliang Guan 已提交
18

19
#include "os.h"
S
Shengliang Guan 已提交
20 21

#include "cJSON.h"
L
Liu Jicong 已提交
22
#include "scheduler.h"
S
Shengliang Guan 已提交
23 24
#include "sync.h"
#include "thash.h"
L
Liu Jicong 已提交
25
#include "tlist.h"
S
Shengliang Guan 已提交
26
#include "tlog.h"
L
Liu Jicong 已提交
27
#include "tmsg.h"
S
Shengliang Guan 已提交
28 29
#include "trpc.h"
#include "ttimer.h"
S
Shengliang Guan 已提交
30

S
Shengliang Guan 已提交
31
#include "mnode.h"
S
Shengliang Guan 已提交
32

S
Shengliang Guan 已提交
33 34 35
#ifdef __cplusplus
extern "C" {
#endif
S
Shengliang Guan 已提交
36

S
Shengliang Guan 已提交
37 38 39
extern int32_t mDebugFlag;

// mnode log function
L
Liu Jicong 已提交
40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75
#define mFatal(...)                                 \
  {                                                 \
    if (mDebugFlag & DEBUG_FATAL) {                 \
      taosPrintLog("MND FATAL ", 255, __VA_ARGS__); \
    }                                               \
  }
#define mError(...)                                 \
  {                                                 \
    if (mDebugFlag & DEBUG_ERROR) {                 \
      taosPrintLog("MND ERROR ", 255, __VA_ARGS__); \
    }                                               \
  }
#define mWarn(...)                                 \
  {                                                \
    if (mDebugFlag & DEBUG_WARN) {                 \
      taosPrintLog("MND WARN ", 255, __VA_ARGS__); \
    }                                              \
  }
#define mInfo(...)                            \
  {                                           \
    if (mDebugFlag & DEBUG_INFO) {            \
      taosPrintLog("MND ", 255, __VA_ARGS__); \
    }                                         \
  }
#define mDebug(...)                                  \
  {                                                  \
    if (mDebugFlag & DEBUG_DEBUG) {                  \
      taosPrintLog("MND ", mDebugFlag, __VA_ARGS__); \
    }                                                \
  }
#define mTrace(...)                                  \
  {                                                  \
    if (mDebugFlag & DEBUG_TRACE) {                  \
      taosPrintLog("MND ", mDebugFlag, __VA_ARGS__); \
    }                                                \
  }
S
Shengliang Guan 已提交
76 77

typedef enum {
78 79 80 81 82 83 84 85
  MND_AUTH_ACCT_START = 0,
  MND_AUTH_ACCT_USER,
  MND_AUTH_ACCT_DNODE,
  MND_AUTH_ACCT_MNODE,
  MND_AUTH_ACCT_DB,
  MND_AUTH_ACCT_TABLE,
  MND_AUTH_ACCT_MAX
} EAuthAcct;
S
Shengliang Guan 已提交
86 87

typedef enum {
88 89 90 91 92 93
  MND_AUTH_OP_START = 0,
  MND_AUTH_OP_CREATE_USER,
  MND_AUTH_OP_ALTER_USER,
  MND_AUTH_OP_DROP_USER,
  MND_AUTH_MAX
} EAuthOp;
S
Shengliang Guan 已提交
94

S
Shengliang Guan 已提交
95
typedef enum {
96
  TRN_STAGE_PREPARE = 0,
S
Shengliang Guan 已提交
97 98
  TRN_STAGE_REDO_LOG = 1,
  TRN_STAGE_REDO_ACTION = 2,
99 100
  TRN_STAGE_COMMIT = 3,
  TRN_STAGE_COMMIT_LOG = 4,
101 102
  TRN_STAGE_UNDO_ACTION = 5,
  TRN_STAGE_UNDO_LOG = 6,
S
Shengliang Guan 已提交
103 104
  TRN_STAGE_ROLLBACK = 7,
  TRN_STAGE_FINISHED = 8
S
Shengliang Guan 已提交
105 106
} ETrnStage;

S
Shengliang Guan 已提交
107
typedef enum {
S
Shengliang Guan 已提交
108 109 110 111
  TRN_TYPE_BASIC_SCOPE = 1000,
  TRN_TYPE_CREATE_USER = 1001,
  TRN_TYPE_ALTER_USER = 1002,
  TRN_TYPE_DROP_USER = 1003,
S
Shengliang Guan 已提交
112 113
  TRN_TYPE_CREATE_FUNC = 1004,
  TRN_TYPE_DROP_FUNC = 1005,
S
Shengliang Guan 已提交
114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140
  TRN_TYPE_CREATE_SNODE = 1006,
  TRN_TYPE_DROP_SNODE = 1007,
  TRN_TYPE_CREATE_QNODE = 1008,
  TRN_TYPE_DROP_QNODE = 1009,
  TRN_TYPE_CREATE_BNODE = 1010,
  TRN_TYPE_DROP_BNODE = 1011,
  TRN_TYPE_CREATE_MNODE = 1012,
  TRN_TYPE_DROP_MNODE = 1013,
  TRN_TYPE_CREATE_TOPIC = 1014,
  TRN_TYPE_DROP_TOPIC = 1015,
  TRN_TYPE_SUBSCRIBE = 1016,
  TRN_TYPE_REBALANCE = 1017,
  TRN_TYPE_BASIC_SCOPE_END,
  TRN_TYPE_GLOBAL_SCOPE = 2000,
  TRN_TYPE_CREATE_DNODE = 2001,
  TRN_TYPE_DROP_DNODE = 2002,
  TRN_TYPE_GLOBAL_SCOPE_END,
  TRN_TYPE_DB_SCOPE = 3000,
  TRN_TYPE_CREATE_DB = 3001,
  TRN_TYPE_ALTER_DB = 3002,
  TRN_TYPE_DROP_DB = 3003,
  TRN_TYPE_CREATE_STB = 3004,
  TRN_TYPE_ALTER_STB = 3005,
  TRN_TYPE_DROP_STB = 3006,
  TRN_TYPE_SPLIT_VGROUP = 3007,
  TRN_TYPE_MERGE_VGROUP = 3018,
  TRN_TYPE_DB_SCOPE_END,
S
Shengliang Guan 已提交
141 142
} ETrnType;

143
typedef enum { TRN_POLICY_ROLLBACK = 0, TRN_POLICY_RETRY = 1 } ETrnPolicy;
S
Shengliang Guan 已提交
144

S
Shengliang Guan 已提交
145 146 147 148 149 150 151 152 153 154 155 156 157 158
typedef enum {
  DND_REASON_ONLINE = 0,
  DND_REASON_STATUS_MSG_TIMEOUT,
  DND_REASON_STATUS_NOT_RECEIVED,
  DND_REASON_VERSION_NOT_MATCH,
  DND_REASON_DNODE_ID_NOT_MATCH,
  DND_REASON_CLUSTER_ID_NOT_MATCH,
  DND_REASON_STATUS_INTERVAL_NOT_MATCH,
  DND_REASON_TIME_ZONE_NOT_MATCH,
  DND_REASON_LOCALE_NOT_MATCH,
  DND_REASON_CHARSET_NOT_MATCH,
  DND_REASON_OTHERS
} EDndReason;

S
Shengliang Guan 已提交
159
typedef struct {
S
Shengliang Guan 已提交
160 161 162
  int32_t    id;
  ETrnStage  stage;
  ETrnPolicy policy;
S
Shengliang Guan 已提交
163
  ETrnType   transType;
S
Shengliang Guan 已提交
164 165
  int32_t    code;
  int32_t    failedTimes;
L
Liu Jicong 已提交
166 167
  void*      rpcHandle;
  void*      rpcAHandle;
S
Shengliang Guan 已提交
168 169
  void*      rpcRsp;
  int32_t    rpcRspLen;
L
Liu Jicong 已提交
170 171 172 173 174
  SArray*    redoLogs;
  SArray*    undoLogs;
  SArray*    commitLogs;
  SArray*    redoActions;
  SArray*    undoActions;
S
Shengliang Guan 已提交
175 176 177
  int64_t    createdTime;
  int64_t    lastExecTime;
  uint64_t   dbUid;
S
Shengliang Guan 已提交
178
  char       dbname[TSDB_DB_FNAME_LEN];
S
Shengliang Guan 已提交
179
  char       lastError[TSDB_TRANS_ERROR_LEN];
S
Shengliang Guan 已提交
180
} STrans;
S
Shengliang Guan 已提交
181

S
Shengliang Guan 已提交
182
typedef struct {
183
  int64_t id;
S
Shengliang Guan 已提交
184
  char    name[TSDB_CLUSTER_ID_LEN];
S
Shengliang Guan 已提交
185 186 187 188
  int64_t createdTime;
  int64_t updateTime;
} SClusterObj;

S
Shengliang Guan 已提交
189
typedef struct {
S
Shengliang Guan 已提交
190 191 192 193
  int32_t    id;
  int64_t    createdTime;
  int64_t    updateTime;
  int64_t    rebootTime;
194
  int64_t    lastAccessTime;
S
Shengliang Guan 已提交
195
  int32_t    accessTimes;
S
Shengliang Guan 已提交
196
  int32_t    numOfVnodes;
S
Shengliang Guan 已提交
197 198
  int32_t    numOfSupportVnodes;
  int32_t    numOfCores;
S
Shengliang Guan 已提交
199 200 201 202
  EDndReason offlineReason;
  uint16_t   port;
  char       fqdn[TSDB_FQDN_LEN];
  char       ep[TSDB_EP_LEN];
S
Shengliang Guan 已提交
203 204
} SDnodeObj;

S
Shengliang Guan 已提交
205
typedef struct {
S
Shengliang Guan 已提交
206
  int32_t    id;
S
Shengliang Guan 已提交
207 208
  int64_t    createdTime;
  int64_t    updateTime;
S
Shengliang Guan 已提交
209
  ESyncState role;
S
Shengliang Guan 已提交
210 211
  int32_t    roleTerm;
  int64_t    roleTime;
L
Liu Jicong 已提交
212
  SDnodeObj* pDnode;
S
Shengliang Guan 已提交
213 214
} SMnodeObj;

S
Shengliang Guan 已提交
215 216 217 218
typedef struct {
  int32_t    id;
  int64_t    createdTime;
  int64_t    updateTime;
L
Liu Jicong 已提交
219
  SDnodeObj* pDnode;
S
Shengliang Guan 已提交
220 221 222 223 224 225
} SQnodeObj;

typedef struct {
  int32_t    id;
  int64_t    createdTime;
  int64_t    updateTime;
L
Liu Jicong 已提交
226
  SDnodeObj* pDnode;
S
Shengliang Guan 已提交
227 228 229 230 231 232
} SSnodeObj;

typedef struct {
  int32_t    id;
  int64_t    createdTime;
  int64_t    updateTime;
L
Liu Jicong 已提交
233
  SDnodeObj* pDnode;
S
Shengliang Guan 已提交
234 235
} SBnodeObj;

S
Shengliang Guan 已提交
236 237 238
typedef struct {
  int32_t maxUsers;
  int32_t maxDbs;
S
Shengliang Guan 已提交
239 240
  int32_t maxStbs;
  int32_t maxTbs;
S
Shengliang Guan 已提交
241 242
  int32_t maxTimeSeries;
  int32_t maxStreams;
S
Shengliang Guan 已提交
243 244 245 246
  int32_t maxFuncs;
  int32_t maxConsumers;
  int32_t maxConns;
  int32_t maxTopics;
S
Shengliang Guan 已提交
247 248
  int64_t maxStorage;   // In unit of GB
  int32_t accessState;  // Configured only by command
S
Shengliang Guan 已提交
249 250 251 252 253 254 255
} SAcctCfg;

typedef struct {
  int32_t numOfUsers;
  int32_t numOfDbs;
  int32_t numOfTimeSeries;
  int32_t numOfStreams;
S
Shengliang Guan 已提交
256 257
  int64_t totalStorage;  // Total storage wrtten from this account
  int64_t compStorage;   // Compressed storage on disk
S
Shengliang Guan 已提交
258 259
} SAcctInfo;

S
Shengliang Guan 已提交
260
typedef struct {
S
Shengliang Guan 已提交
261 262 263 264
  char      acct[TSDB_USER_LEN];
  int64_t   createdTime;
  int64_t   updateTime;
  int32_t   acctId;
S
Shengliang Guan 已提交
265
  int32_t   status;
S
Shengliang Guan 已提交
266 267 268 269
  SAcctCfg  cfg;
  SAcctInfo info;
} SAcctObj;

S
Shengliang Guan 已提交
270
typedef struct {
S
Shengliang Guan 已提交
271
  char      user[TSDB_USER_LEN];
272
  char      pass[TSDB_PASSWORD_LEN];
S
Shengliang Guan 已提交
273 274 275
  char      acct[TSDB_USER_LEN];
  int64_t   createdTime;
  int64_t   updateTime;
276
  int8_t    superUser;
S
Shengliang Guan 已提交
277
  int32_t   acctId;
278 279
  SHashObj* readDbs;
  SHashObj* writeDbs;
S
Shengliang Guan 已提交
280 281 282
} SUserObj;

typedef struct {
283
  int32_t numOfVgroups;
S
Shengliang Guan 已提交
284 285 286 287 288 289
  int32_t cacheBlockSize;
  int32_t totalBlocks;
  int32_t daysPerFile;
  int32_t daysToKeep0;
  int32_t daysToKeep1;
  int32_t daysToKeep2;
S
Shengliang Guan 已提交
290 291
  int32_t minRows;
  int32_t maxRows;
S
Shengliang Guan 已提交
292 293
  int32_t commitTime;
  int32_t fsyncPeriod;
S
Shengliang Guan 已提交
294
  int8_t  walLevel;
S
Shengliang Guan 已提交
295 296 297 298 299 300 301 302
  int8_t  precision;
  int8_t  compression;
  int8_t  replications;
  int8_t  quorum;
  int8_t  update;
  int8_t  cacheLastRow;
} SDbCfg;

S
Shengliang Guan 已提交
303
typedef struct {
L
Liu Jicong 已提交
304 305
  char     name[TSDB_DB_FNAME_LEN];
  char     acct[TSDB_USER_LEN];
S
Shengliang Guan 已提交
306
  char     createUser[TSDB_USER_LEN];
L
Liu Jicong 已提交
307 308
  int64_t  createdTime;
  int64_t  updateTime;
H
Haojun Liao 已提交
309
  uint64_t uid;
L
Liu Jicong 已提交
310 311 312 313
  int32_t  cfgVersion;
  int32_t  vgVersion;
  int8_t   hashMethod;  // default is 1
  SDbCfg   cfg;
S
Shengliang Guan 已提交
314 315 316 317
} SDbObj;

typedef struct {
  int32_t    dnodeId;
S
Shengliang Guan 已提交
318
  ESyncState role;
S
Shengliang Guan 已提交
319 320
} SVnodeGid;

S
Shengliang Guan 已提交
321
typedef struct {
S
Shengliang Guan 已提交
322
  int32_t   vgId;
S
Shengliang Guan 已提交
323 324
  int64_t   createdTime;
  int64_t   updateTime;
S
Shengliang Guan 已提交
325
  int32_t   version;
S
Shengliang Guan 已提交
326 327
  uint32_t  hashBegin;
  uint32_t  hashEnd;
328
  char      dbName[TSDB_DB_FNAME_LEN];
S
Shengliang Guan 已提交
329
  int64_t   dbUid;
S
Shengliang Guan 已提交
330 331
  int64_t   numOfTables;
  int64_t   numOfTimeSeries;
S
Shengliang Guan 已提交
332 333 334
  int64_t   totalStorage;
  int64_t   compStorage;
  int64_t   pointsWritten;
S
Shengliang Guan 已提交
335 336 337
  int8_t    compact;
  int8_t    replica;
  SVnodeGid vnodeGid[TSDB_MAX_REPLICA];
S
Shengliang Guan 已提交
338 339
} SVgObj;

S
Shengliang Guan 已提交
340
typedef struct {
S
Shengliang Guan 已提交
341
  char     name[TSDB_TABLE_FNAME_LEN];
342
  char     db[TSDB_DB_FNAME_LEN];
343 344
  int64_t  createdTime;
  int64_t  updateTime;
S
Shengliang Guan 已提交
345
  uint64_t uid;
S
Shengliang Guan 已提交
346
  uint64_t dbUid;
S
Shengliang Guan 已提交
347
  int32_t  version;
S
Shengliang Guan 已提交
348
  int32_t  nextColId;
S
Shengliang Guan 已提交
349 350
  int32_t  numOfColumns;
  int32_t  numOfTags;
S
Shengliang Guan 已提交
351
  SSchema* pColumns;
S
Shengliang Guan 已提交
352
  SSchema* pTags;
S
Shengliang Guan 已提交
353
  SRWLatch lock;
S
Shengliang Guan 已提交
354
  char     comment[TSDB_STB_COMMENT_LEN];
S
Shengliang Guan 已提交
355
} SStbObj;
S
Shengliang Guan 已提交
356

S
Shengliang Guan 已提交
357
typedef struct {
S
Shengliang Guan 已提交
358 359
  char    name[TSDB_FUNC_NAME_LEN];
  int64_t createdTime;
S
Shengliang Guan 已提交
360 361 362 363 364 365
  int8_t  funcType;
  int8_t  scriptType;
  int8_t  align;
  int8_t  outputType;
  int32_t outputLen;
  int32_t bufSize;
S
Shengliang 已提交
366
  int64_t signature;
S
Shengliang Guan 已提交
367 368
  int32_t commentSize;
  int32_t codeSize;
L
Liu Jicong 已提交
369 370
  char*   pComment;
  char*   pCode;
S
Shengliang Guan 已提交
371
  char    pData[];
S
Shengliang Guan 已提交
372 373
} SFuncObj;

S
Shengliang Guan 已提交
374
typedef struct {
375
  int64_t id;
S
Shengliang Guan 已提交
376 377 378 379 380 381 382
  int8_t  type;
  int8_t  replica;
  int16_t numOfColumns;
  int32_t rowSize;
  int32_t numOfRows;
  int32_t numOfReads;
  int32_t payloadLen;
L
Liu Jicong 已提交
383 384
  void*   pIter;
  SMnode* pMnode;
385
  char    db[TSDB_DB_FNAME_LEN];
S
Shengliang Guan 已提交
386 387 388
  int16_t offset[TSDB_MAX_COLUMNS];
  int32_t bytes[TSDB_MAX_COLUMNS];
  char    payload[];
S
Shengliang Guan 已提交
389 390
} SShowObj;

L
Liu Jicong 已提交
391 392 393 394 395 396 397
typedef struct {
  int32_t vgId;  // -1 for unassigned
  int32_t status;
  SEpSet  epSet;
  int64_t oldConsumerId;
  int64_t consumerId;  // -1 for unassigned
  char*   qmsg;
L
Liu Jicong 已提交
398 399
} SMqConsumerEp;

L
Liu Jicong 已提交
400
static FORCE_INLINE int32_t tEncodeSMqConsumerEp(void** buf, const SMqConsumerEp* pConsumerEp) {
L
Liu Jicong 已提交
401 402
  int32_t tlen = 0;
  tlen += taosEncodeFixedI32(buf, pConsumerEp->vgId);
L
Liu Jicong 已提交
403
  tlen += taosEncodeFixedI32(buf, pConsumerEp->status);
404
  tlen += taosEncodeSEpSet(buf, &pConsumerEp->epSet);
L
Liu Jicong 已提交
405
  tlen += taosEncodeFixedI64(buf, pConsumerEp->oldConsumerId);
L
Liu Jicong 已提交
406
  tlen += taosEncodeFixedI64(buf, pConsumerEp->consumerId);
L
Liu Jicong 已提交
407
  tlen += taosEncodeString(buf, pConsumerEp->qmsg);
L
Liu Jicong 已提交
408 409 410 411 412
  return tlen;
}

static FORCE_INLINE void* tDecodeSMqConsumerEp(void** buf, SMqConsumerEp* pConsumerEp) {
  buf = taosDecodeFixedI32(buf, &pConsumerEp->vgId);
L
Liu Jicong 已提交
413
  buf = taosDecodeFixedI32(buf, &pConsumerEp->status);
414
  buf = taosDecodeSEpSet(buf, &pConsumerEp->epSet);
L
Liu Jicong 已提交
415
  buf = taosDecodeFixedI64(buf, &pConsumerEp->oldConsumerId);
L
Liu Jicong 已提交
416
  buf = taosDecodeFixedI64(buf, &pConsumerEp->consumerId);
L
Liu Jicong 已提交
417
  buf = taosDecodeString(buf, &pConsumerEp->qmsg);
L
Liu Jicong 已提交
418 419 420
  return buf;
}

L
Liu Jicong 已提交
421 422 423 424 425 426
static FORCE_INLINE void tDeleteSMqConsumerEp(SMqConsumerEp* pConsumerEp) {
  if (pConsumerEp) {
    tfree(pConsumerEp->qmsg);
  }
}

L
Liu Jicong 已提交
427 428 429 430
typedef struct {
  int64_t consumerId;
  SArray* vgInfo;  // SArray<SMqConsumerEp>
} SMqSubConsumer;
L
Liu Jicong 已提交
431

L
Liu Jicong 已提交
432 433 434 435
static FORCE_INLINE int32_t tEncodeSMqSubConsumer(void** buf, const SMqSubConsumer* pConsumer) {
  int32_t tlen = 0;
  tlen += taosEncodeFixedI64(buf, pConsumer->consumerId);
  int32_t sz = taosArrayGetSize(pConsumer->vgInfo);
L
Liu Jicong 已提交
436
  tlen += taosEncodeFixedI32(buf, sz);
L
Liu Jicong 已提交
437 438 439
  for (int32_t i = 0; i < sz; i++) {
    SMqConsumerEp* pCEp = taosArrayGet(pConsumer->vgInfo, i);
    tlen += tEncodeSMqConsumerEp(buf, pCEp);
L
Liu Jicong 已提交
440
  }
L
Liu Jicong 已提交
441 442
  return tlen;
}
L
Liu Jicong 已提交
443

L
Liu Jicong 已提交
444 445 446 447 448 449 450 451 452
static FORCE_INLINE void* tDecodeSMqSubConsumer(void** buf, SMqSubConsumer* pConsumer) {
  int32_t sz;
  buf = taosDecodeFixedI64(buf, &pConsumer->consumerId);
  buf = taosDecodeFixedI32(buf, &sz);
  pConsumer->vgInfo = taosArrayInit(sz, sizeof(SMqConsumerEp));
  for (int32_t i = 0; i < sz; i++) {
    SMqConsumerEp consumerEp;
    buf = tDecodeSMqConsumerEp(buf, &consumerEp);
    taosArrayPush(pConsumer->vgInfo, &consumerEp);
L
Liu Jicong 已提交
453
  }
L
Liu Jicong 已提交
454 455 456 457 458 459 460
  return buf;
}

static FORCE_INLINE void tDeleteSMqSubConsumer(SMqSubConsumer* pSubConsumer) {
  if (pSubConsumer->vgInfo) {
    taosArrayDestroyEx(pSubConsumer->vgInfo, (void (*)(void*))tDeleteSMqConsumerEp);
    pSubConsumer->vgInfo = NULL;
L
Liu Jicong 已提交
461
  }
L
Liu Jicong 已提交
462 463 464 465 466 467 468
}

typedef struct {
  char    key[TSDB_SUBSCRIBE_KEY_LEN];
  int32_t status;
  int32_t vgNum;
  SArray* consumers;     // SArray<SMqSubConsumer>
L
Liu Jicong 已提交
469
  SArray* lostConsumers; // SArray<SMqSubConsumer>
L
Liu Jicong 已提交
470 471 472 473 474 475
  SArray* unassignedVg;  // SArray<SMqConsumerEp>
} SMqSubscribeObj;

static FORCE_INLINE SMqSubscribeObj* tNewSubscribeObj() {
  SMqSubscribeObj* pSub = calloc(1, sizeof(SMqSubscribeObj));
  if (pSub == NULL) {
L
Liu Jicong 已提交
476 477
    return NULL;
  }
L
Liu Jicong 已提交
478

L
Liu Jicong 已提交
479 480 481
  pSub->consumers = taosArrayInit(0, sizeof(SMqSubConsumer));
  if (pSub->consumers == NULL) {
    goto _err;
L
Liu Jicong 已提交
482
  }
L
Liu Jicong 已提交
483 484 485 486 487 488

  pSub->lostConsumers = taosArrayInit(0, sizeof(SMqSubConsumer));
  if (pSub->lostConsumers == NULL) {
    goto _err;
  }

L
Liu Jicong 已提交
489
  pSub->unassignedVg = taosArrayInit(0, sizeof(SMqConsumerEp));
L
Liu Jicong 已提交
490
  if (pSub->unassignedVg == NULL) {
L
Liu Jicong 已提交
491
    goto _err;
L
Liu Jicong 已提交
492
  }
L
Liu Jicong 已提交
493 494 495 496 497

  pSub->key[0] = 0;
  pSub->vgNum = 0;
  pSub->status = 0;

L
Liu Jicong 已提交
498
  return pSub;
L
Liu Jicong 已提交
499 500 501

_err:
  tfree(pSub->consumers);
L
Liu Jicong 已提交
502 503
  tfree(pSub->lostConsumers);
  tfree(pSub->unassignedVg);
L
Liu Jicong 已提交
504 505
  tfree(pSub);
  return NULL;
L
Liu Jicong 已提交
506 507 508 509 510
}

static FORCE_INLINE int32_t tEncodeSubscribeObj(void** buf, const SMqSubscribeObj* pSub) {
  int32_t tlen = 0;
  tlen += taosEncodeString(buf, pSub->key);
L
Liu Jicong 已提交
511 512
  tlen += taosEncodeFixedI32(buf, pSub->vgNum);
  tlen += taosEncodeFixedI32(buf, pSub->status);
L
Liu Jicong 已提交
513 514
  int32_t sz;

L
Liu Jicong 已提交
515
  sz = taosArrayGetSize(pSub->consumers);
L
Liu Jicong 已提交
516 517
  tlen += taosEncodeFixedI32(buf, sz);
  for (int32_t i = 0; i < sz; i++) {
L
Liu Jicong 已提交
518 519
    SMqSubConsumer* pSubConsumer = taosArrayGet(pSub->consumers, i);
    tlen += tEncodeSMqSubConsumer(buf, pSubConsumer);
L
Liu Jicong 已提交
520 521
  }

L
Liu Jicong 已提交
522 523 524 525 526 527 528
  sz = taosArrayGetSize(pSub->lostConsumers);
  tlen += taosEncodeFixedI32(buf, sz);
  for (int32_t i = 0; i < sz; i++) {
    SMqSubConsumer* pSubConsumer = taosArrayGet(pSub->lostConsumers, i);
    tlen += tEncodeSMqSubConsumer(buf, pSubConsumer);
  }

L
Liu Jicong 已提交
529 530 531 532 533 534 535 536 537 538 539 540
  sz = taosArrayGetSize(pSub->unassignedVg);
  tlen += taosEncodeFixedI32(buf, sz);
  for (int32_t i = 0; i < sz; i++) {
    SMqConsumerEp* pCEp = taosArrayGet(pSub->unassignedVg, i);
    tlen += tEncodeSMqConsumerEp(buf, pCEp);
  }

  return tlen;
}

static FORCE_INLINE void* tDecodeSubscribeObj(void* buf, SMqSubscribeObj* pSub) {
  buf = taosDecodeStringTo(buf, pSub->key);
L
Liu Jicong 已提交
541 542
  buf = taosDecodeFixedI32(buf, &pSub->vgNum);
  buf = taosDecodeFixedI32(buf, &pSub->status);
L
Liu Jicong 已提交
543 544 545 546

  int32_t sz;

  buf = taosDecodeFixedI32(buf, &sz);
L
Liu Jicong 已提交
547 548
  pSub->consumers = taosArrayInit(sz, sizeof(SMqSubConsumer));
  if (pSub->consumers == NULL) {
L
Liu Jicong 已提交
549 550 551
    return NULL;
  }
  for (int32_t i = 0; i < sz; i++) {
L
Liu Jicong 已提交
552 553 554
    SMqSubConsumer subConsumer = {0};
    buf = tDecodeSMqSubConsumer(buf, &subConsumer);
    taosArrayPush(pSub->consumers, &subConsumer);
L
Liu Jicong 已提交
555 556
  }

L
Liu Jicong 已提交
557 558 559 560 561 562 563 564 565 566 567
  buf = taosDecodeFixedI32(buf, &sz);
  pSub->lostConsumers = taosArrayInit(sz, sizeof(SMqSubConsumer));
  if (pSub->lostConsumers == NULL) {
    return NULL;
  }
  for (int32_t i = 0; i < sz; i++) {
    SMqSubConsumer subConsumer = {0};
    buf = tDecodeSMqSubConsumer(buf, &subConsumer);
    taosArrayPush(pSub->lostConsumers, &subConsumer);
  }

L
Liu Jicong 已提交
568 569 570 571 572 573
  buf = taosDecodeFixedI32(buf, &sz);
  pSub->unassignedVg = taosArrayInit(sz, sizeof(SMqConsumerEp));
  if (pSub->unassignedVg == NULL) {
    return NULL;
  }
  for (int32_t i = 0; i < sz; i++) {
L
Liu Jicong 已提交
574 575 576
    SMqConsumerEp consumerEp = {0};
    buf = tDecodeSMqConsumerEp(buf, &consumerEp);
    taosArrayPush(pSub->unassignedVg, &consumerEp);
L
Liu Jicong 已提交
577 578 579
  }
  return buf;
}
L
Liu Jicong 已提交
580

L
Liu Jicong 已提交
581
static FORCE_INLINE void tDeleteSMqSubscribeObj(SMqSubscribeObj* pSub) {
L
Liu Jicong 已提交
582 583 584 585
  if (pSub->consumers) {
    taosArrayDestroyEx(pSub->consumers, (void (*)(void*))tDeleteSMqSubConsumer);
    //taosArrayDestroy(pSub->consumers);
    pSub->consumers = NULL;
L
Liu Jicong 已提交
586
  }
L
Liu Jicong 已提交
587

L
Liu Jicong 已提交
588
  if (pSub->unassignedVg) {
L
Liu Jicong 已提交
589 590
    taosArrayDestroyEx(pSub->unassignedVg, (void (*)(void*))tDeleteSMqConsumerEp);
    //taosArrayDestroy(pSub->unassignedVg);
L
Liu Jicong 已提交
591 592
    pSub->unassignedVg = NULL;
  }
L
Liu Jicong 已提交
593 594
}

L
Liu Jicong 已提交
595
typedef struct {
L
Liu Jicong 已提交
596 597 598 599
  char     name[TSDB_TOPIC_FNAME_LEN];
  char     db[TSDB_DB_FNAME_LEN];
  int64_t  createTime;
  int64_t  updateTime;
L
Liu Jicong 已提交
600
  int64_t  uid;
L
Liu Jicong 已提交
601
  int64_t  dbUid;
L
Liu Jicong 已提交
602 603 604 605 606 607
  int32_t  version;
  SRWLatch lock;
  int32_t  sqlLen;
  char*    sql;
  char*    logicalPlan;
  char*    physicalPlan;
L
Liu Jicong 已提交
608 609
} SMqTopicObj;

L
Liu Jicong 已提交
610
typedef struct {
L
Liu Jicong 已提交
611
  int64_t  consumerId;
L
Liu Jicong 已提交
612
  int64_t  connId;
L
Liu Jicong 已提交
613
  SRWLatch lock;
L
Liu Jicong 已提交
614
  char     cgroup[TSDB_CONSUMER_GROUP_LEN];
L
Liu Jicong 已提交
615 616
  SArray*  currentTopics;  // SArray<char*>
  SArray*  recentRemovedTopics;   // SArray<char*>
L
Liu Jicong 已提交
617
  int32_t  epoch;
L
Liu Jicong 已提交
618
  // stat
L
Liu Jicong 已提交
619 620 621 622 623 624 625
  int64_t pollCnt;
  // status
  int32_t status;
  // heartbeat from the consumer reset hbStatus to 0
  // each checkConsumerAlive msg add hbStatus by 1
  // if checkConsumerAlive > CONSUMER_REBALANCE_CNT, mask to lost
  int32_t hbStatus;
L
Liu Jicong 已提交
626 627
} SMqConsumerObj;

L
Liu Jicong 已提交
628
static FORCE_INLINE int32_t tEncodeSMqConsumerObj(void** buf, const SMqConsumerObj* pConsumer) {
L
Liu Jicong 已提交
629
  int32_t sz;
L
Liu Jicong 已提交
630 631
  int32_t tlen = 0;
  tlen += taosEncodeFixedI64(buf, pConsumer->consumerId);
L
Liu Jicong 已提交
632
  tlen += taosEncodeFixedI64(buf, pConsumer->connId);
L
Liu Jicong 已提交
633
  tlen += taosEncodeFixedI32(buf, pConsumer->epoch);
L
Liu Jicong 已提交
634
  tlen += taosEncodeFixedI64(buf, pConsumer->pollCnt);
L
Liu Jicong 已提交
635
  tlen += taosEncodeFixedI32(buf, pConsumer->status);
L
Liu Jicong 已提交
636
  tlen += taosEncodeString(buf, pConsumer->cgroup);
L
Liu Jicong 已提交
637 638 639 640 641 642 643 644 645

  sz = taosArrayGetSize(pConsumer->currentTopics);
  tlen += taosEncodeFixedI32(buf, sz);
  for (int32_t i = 0; i < sz; i++) {
    char* topic = taosArrayGetP(pConsumer->currentTopics, i);
    tlen += taosEncodeString(buf, topic);
  }

  sz = taosArrayGetSize(pConsumer->recentRemovedTopics);
L
Liu Jicong 已提交
646 647
  tlen += taosEncodeFixedI32(buf, sz);
  for (int32_t i = 0; i < sz; i++) {
L
Liu Jicong 已提交
648
    char* topic = taosArrayGetP(pConsumer->recentRemovedTopics, i);
L
Liu Jicong 已提交
649
    tlen += taosEncodeString(buf, topic);
L
Liu Jicong 已提交
650 651 652 653 654
  }
  return tlen;
}

static FORCE_INLINE void* tDecodeSMqConsumerObj(void* buf, SMqConsumerObj* pConsumer) {
L
Liu Jicong 已提交
655
  int32_t sz;
L
Liu Jicong 已提交
656
  buf = taosDecodeFixedI64(buf, &pConsumer->consumerId);
L
Liu Jicong 已提交
657
  buf = taosDecodeFixedI64(buf, &pConsumer->connId);
L
Liu Jicong 已提交
658
  buf = taosDecodeFixedI32(buf, &pConsumer->epoch);
L
Liu Jicong 已提交
659
  buf = taosDecodeFixedI64(buf, &pConsumer->pollCnt);
L
Liu Jicong 已提交
660
  buf = taosDecodeFixedI32(buf, &pConsumer->status);
L
Liu Jicong 已提交
661
  buf = taosDecodeStringTo(buf, pConsumer->cgroup);
L
Liu Jicong 已提交
662 663 664 665 666 667 668 669 670

  buf = taosDecodeFixedI32(buf, &sz);
  pConsumer->currentTopics = taosArrayInit(sz, sizeof(SMqConsumerObj));
  for (int32_t i = 0; i < sz; i++) {
    char* topic;
    buf = taosDecodeString(buf, &topic);
    taosArrayPush(pConsumer->currentTopics, &topic);
  }

L
Liu Jicong 已提交
671
  buf = taosDecodeFixedI32(buf, &sz);
L
Liu Jicong 已提交
672
  pConsumer->recentRemovedTopics = taosArrayInit(sz, sizeof(SMqConsumerObj));
L
Liu Jicong 已提交
673
  for (int32_t i = 0; i < sz; i++) {
L
Liu Jicong 已提交
674 675
    char* topic;
    buf = taosDecodeString(buf, &topic);
L
Liu Jicong 已提交
676
    taosArrayPush(pConsumer->recentRemovedTopics, &topic);
L
Liu Jicong 已提交
677 678 679 680
  }
  return buf;
}

S
Shengliang Guan 已提交
681
typedef struct SMnodeMsg {
S
Shengliang Guan 已提交
682
  char    user[TSDB_USER_LEN];
683
  char    db[TSDB_DB_FNAME_LEN];
S
Shengliang Guan 已提交
684
  int32_t acctId;
L
Liu Jicong 已提交
685
  SMnode* pMnode;
S
Shengliang Guan 已提交
686 687 688
  int64_t createdTime;
  SRpcMsg rpcMsg;
  int32_t contLen;
L
Liu Jicong 已提交
689
  void*   pCont;
S
Shengliang Guan 已提交
690
} SMnodeMsg;
S
Shengliang Guan 已提交
691

S
Shengliang Guan 已提交
692 693 694
#ifdef __cplusplus
}
#endif
S
Shengliang Guan 已提交
695

S
Shengliang Guan 已提交
696
#endif /*_TD_MND_DEF_H_*/