mndDef.h 19.5 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
#include "trpc.h"
L
Liu Jicong 已提交
29
#include "tstream.h"
S
Shengliang Guan 已提交
30
#include "ttimer.h"
S
Shengliang Guan 已提交
31

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

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

S
Shengliang Guan 已提交
38
typedef enum {
39 40 41 42 43 44 45 46
  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 已提交
47 48

typedef enum {
49 50 51 52 53 54
  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 已提交
55

S
Shengliang Guan 已提交
56
typedef enum {
57
  TRN_STAGE_PREPARE = 0,
S
Shengliang Guan 已提交
58 59
  TRN_STAGE_REDO_LOG = 1,
  TRN_STAGE_REDO_ACTION = 2,
60 61
  TRN_STAGE_COMMIT = 3,
  TRN_STAGE_COMMIT_LOG = 4,
62 63
  TRN_STAGE_UNDO_ACTION = 5,
  TRN_STAGE_UNDO_LOG = 6,
S
Shengliang Guan 已提交
64 65
  TRN_STAGE_ROLLBACK = 7,
  TRN_STAGE_FINISHED = 8
S
Shengliang Guan 已提交
66 67
} ETrnStage;

S
Shengliang Guan 已提交
68
typedef enum {
S
Shengliang Guan 已提交
69 70 71 72
  TRN_TYPE_BASIC_SCOPE = 1000,
  TRN_TYPE_CREATE_USER = 1001,
  TRN_TYPE_ALTER_USER = 1002,
  TRN_TYPE_DROP_USER = 1003,
S
Shengliang Guan 已提交
73 74
  TRN_TYPE_CREATE_FUNC = 1004,
  TRN_TYPE_DROP_FUNC = 1005,
S
Shengliang Guan 已提交
75 76 77 78 79 80 81 82 83 84 85 86
  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,
L
Liu Jicong 已提交
87
  TRN_TYPE_COMMIT_OFFSET = 1018,
L
Liu Jicong 已提交
88
  TRN_TYPE_CREATE_STREAM = 1019,
L
Liu Jicong 已提交
89 90
  TRN_TYPE_DROP_STREAM = 1020,
  TRN_TYPE_ALTER_STREAM = 1021,
S
Shengliang Guan 已提交
91 92 93 94 95 96 97 98 99
  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,
S
Shengliang Guan 已提交
100 101
  TRN_TYPE_SPLIT_VGROUP = 3004,
  TRN_TYPE_MERGE_VGROUP = 3015,
S
Shengliang Guan 已提交
102
  TRN_TYPE_DB_SCOPE_END,
S
Shengliang Guan 已提交
103 104 105 106
  TRN_TYPE_STB_SCOPE = 4000,
  TRN_TYPE_CREATE_STB = 4001,
  TRN_TYPE_ALTER_STB = 4002,
  TRN_TYPE_DROP_STB = 4003,
S
Shengliang Guan 已提交
107 108
  TRN_TYPE_CREATE_SMA = 4004,
  TRN_TYPE_DROP_SMA = 4005,
S
Shengliang Guan 已提交
109
  TRN_TYPE_STB_SCOPE_END,
S
Shengliang Guan 已提交
110 111
} ETrnType;

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

S
Shengliang Guan 已提交
114 115 116 117 118 119 120 121 122 123 124 125 126 127
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 已提交
128
typedef struct {
S
Shengliang Guan 已提交
129 130 131
  int32_t    id;
  ETrnStage  stage;
  ETrnPolicy policy;
S
Shengliang Guan 已提交
132
  ETrnType   transType;
S
Shengliang Guan 已提交
133 134
  int32_t    code;
  int32_t    failedTimes;
L
Liu Jicong 已提交
135 136
  void*      rpcHandle;
  void*      rpcAHandle;
S
Shengliang Guan 已提交
137 138
  void*      rpcRsp;
  int32_t    rpcRspLen;
L
Liu Jicong 已提交
139 140 141 142 143
  SArray*    redoLogs;
  SArray*    undoLogs;
  SArray*    commitLogs;
  SArray*    redoActions;
  SArray*    undoActions;
S
Shengliang Guan 已提交
144 145
  int64_t    createdTime;
  int64_t    lastExecTime;
L
Liu Jicong 已提交
146
  int64_t    dbUid;
S
Shengliang Guan 已提交
147
  char       dbname[TSDB_DB_FNAME_LEN];
S
Shengliang Guan 已提交
148
  char       lastError[TSDB_TRANS_ERROR_LEN];
S
Shengliang Guan 已提交
149
} STrans;
S
Shengliang Guan 已提交
150

S
Shengliang Guan 已提交
151
typedef struct {
152
  int64_t id;
S
Shengliang Guan 已提交
153
  char    name[TSDB_CLUSTER_ID_LEN];
S
Shengliang Guan 已提交
154 155 156 157
  int64_t createdTime;
  int64_t updateTime;
} SClusterObj;

S
Shengliang Guan 已提交
158
typedef struct {
S
Shengliang Guan 已提交
159 160 161 162
  int32_t    id;
  int64_t    createdTime;
  int64_t    updateTime;
  int64_t    rebootTime;
163
  int64_t    lastAccessTime;
S
Shengliang Guan 已提交
164
  int32_t    accessTimes;
S
Shengliang Guan 已提交
165
  int32_t    numOfVnodes;
S
Shengliang Guan 已提交
166 167
  int32_t    numOfSupportVnodes;
  int32_t    numOfCores;
S
Shengliang Guan 已提交
168 169 170 171
  EDndReason offlineReason;
  uint16_t   port;
  char       fqdn[TSDB_FQDN_LEN];
  char       ep[TSDB_EP_LEN];
S
Shengliang Guan 已提交
172 173
} SDnodeObj;

S
Shengliang Guan 已提交
174
typedef struct {
S
Shengliang Guan 已提交
175
  int32_t    id;
S
Shengliang Guan 已提交
176 177
  int64_t    createdTime;
  int64_t    updateTime;
S
Shengliang Guan 已提交
178
  ESyncState role;
S
Shengliang Guan 已提交
179 180
  int32_t    roleTerm;
  int64_t    roleTime;
L
Liu Jicong 已提交
181
  SDnodeObj* pDnode;
S
Shengliang Guan 已提交
182 183
} SMnodeObj;

S
Shengliang Guan 已提交
184 185 186 187
typedef struct {
  int32_t    id;
  int64_t    createdTime;
  int64_t    updateTime;
L
Liu Jicong 已提交
188
  SDnodeObj* pDnode;
S
Shengliang Guan 已提交
189 190 191 192 193 194
} SQnodeObj;

typedef struct {
  int32_t    id;
  int64_t    createdTime;
  int64_t    updateTime;
L
Liu Jicong 已提交
195
  SDnodeObj* pDnode;
S
Shengliang Guan 已提交
196 197 198 199 200 201
} SSnodeObj;

typedef struct {
  int32_t    id;
  int64_t    createdTime;
  int64_t    updateTime;
L
Liu Jicong 已提交
202
  SDnodeObj* pDnode;
S
Shengliang Guan 已提交
203 204
} SBnodeObj;

S
Shengliang Guan 已提交
205 206 207
typedef struct {
  int32_t maxUsers;
  int32_t maxDbs;
S
Shengliang Guan 已提交
208 209
  int32_t maxStbs;
  int32_t maxTbs;
S
Shengliang Guan 已提交
210 211
  int32_t maxTimeSeries;
  int32_t maxStreams;
S
Shengliang Guan 已提交
212 213 214 215
  int32_t maxFuncs;
  int32_t maxConsumers;
  int32_t maxConns;
  int32_t maxTopics;
S
Shengliang Guan 已提交
216 217
  int64_t maxStorage;   // In unit of GB
  int32_t accessState;  // Configured only by command
S
Shengliang Guan 已提交
218 219 220 221 222 223 224
} SAcctCfg;

typedef struct {
  int32_t numOfUsers;
  int32_t numOfDbs;
  int32_t numOfTimeSeries;
  int32_t numOfStreams;
S
Shengliang Guan 已提交
225 226
  int64_t totalStorage;  // Total storage wrtten from this account
  int64_t compStorage;   // Compressed storage on disk
S
Shengliang Guan 已提交
227 228
} SAcctInfo;

S
Shengliang Guan 已提交
229
typedef struct {
S
Shengliang Guan 已提交
230 231 232 233
  char      acct[TSDB_USER_LEN];
  int64_t   createdTime;
  int64_t   updateTime;
  int32_t   acctId;
S
Shengliang Guan 已提交
234
  int32_t   status;
S
Shengliang Guan 已提交
235 236 237 238
  SAcctCfg  cfg;
  SAcctInfo info;
} SAcctObj;

S
Shengliang Guan 已提交
239
typedef struct {
S
Shengliang Guan 已提交
240
  char      user[TSDB_USER_LEN];
241
  char      pass[TSDB_PASSWORD_LEN];
S
Shengliang Guan 已提交
242 243 244
  char      acct[TSDB_USER_LEN];
  int64_t   createdTime;
  int64_t   updateTime;
245
  int8_t    superUser;
S
Shengliang Guan 已提交
246
  int32_t   acctId;
247 248
  SHashObj* readDbs;
  SHashObj* writeDbs;
S
Shengliang Guan 已提交
249 250 251
} SUserObj;

typedef struct {
252
  int32_t numOfVgroups;
S
Shengliang Guan 已提交
253 254 255 256 257 258
  int32_t cacheBlockSize;
  int32_t totalBlocks;
  int32_t daysPerFile;
  int32_t daysToKeep0;
  int32_t daysToKeep1;
  int32_t daysToKeep2;
S
Shengliang Guan 已提交
259 260
  int32_t minRows;
  int32_t maxRows;
S
Shengliang Guan 已提交
261 262
  int32_t commitTime;
  int32_t fsyncPeriod;
D
dapan1121 已提交
263
  int32_t ttl;
S
Shengliang Guan 已提交
264
  int8_t  walLevel;
S
Shengliang Guan 已提交
265 266 267 268 269 270
  int8_t  precision;
  int8_t  compression;
  int8_t  replications;
  int8_t  quorum;
  int8_t  update;
  int8_t  cacheLastRow;
L
Liu Jicong 已提交
271
  int8_t  streamMode;
D
dapan1121 已提交
272
  int8_t  singleSTable;
S
sma  
Shengliang Guan 已提交
273 274
  int32_t numOfRetensions;
  SArray* pRetensions;
S
Shengliang Guan 已提交
275 276
} SDbCfg;

S
Shengliang Guan 已提交
277
typedef struct {
L
Liu Jicong 已提交
278 279 280 281 282 283 284 285 286 287
  char    name[TSDB_DB_FNAME_LEN];
  char    acct[TSDB_USER_LEN];
  char    createUser[TSDB_USER_LEN];
  int64_t createdTime;
  int64_t updateTime;
  int64_t uid;
  int32_t cfgVersion;
  int32_t vgVersion;
  int8_t  hashMethod;  // default is 1
  SDbCfg  cfg;
S
Shengliang Guan 已提交
288 289 290 291
} SDbObj;

typedef struct {
  int32_t    dnodeId;
S
Shengliang Guan 已提交
292
  ESyncState role;
S
Shengliang Guan 已提交
293 294
} SVnodeGid;

S
Shengliang Guan 已提交
295
typedef struct {
S
Shengliang Guan 已提交
296
  int32_t   vgId;
S
Shengliang Guan 已提交
297 298
  int64_t   createdTime;
  int64_t   updateTime;
S
Shengliang Guan 已提交
299
  int32_t   version;
S
Shengliang Guan 已提交
300 301
  uint32_t  hashBegin;
  uint32_t  hashEnd;
302
  char      dbName[TSDB_DB_FNAME_LEN];
S
Shengliang Guan 已提交
303
  int64_t   dbUid;
S
Shengliang Guan 已提交
304 305
  int64_t   numOfTables;
  int64_t   numOfTimeSeries;
S
Shengliang Guan 已提交
306 307 308
  int64_t   totalStorage;
  int64_t   compStorage;
  int64_t   pointsWritten;
S
Shengliang Guan 已提交
309 310
  int8_t    compact;
  int8_t    replica;
L
Liu Jicong 已提交
311
  int8_t    streamMode;
S
Shengliang Guan 已提交
312
  SVnodeGid vnodeGid[TSDB_MAX_REPLICA];
S
Shengliang Guan 已提交
313 314
} SVgObj;

S
Shengliang Guan 已提交
315
typedef struct {
S
Shengliang Guan 已提交
316
  char    name[TSDB_TABLE_FNAME_LEN];
S
sma  
Shengliang Guan 已提交
317 318
  char    stb[TSDB_TABLE_FNAME_LEN];
  char    db[TSDB_DB_FNAME_LEN];
S
Shengliang Guan 已提交
319 320 321
  int64_t createdTime;
  int64_t uid;
  int64_t stbUid;
S
sma  
Shengliang Guan 已提交
322
  int64_t dbUid;
S
Shengliang Guan 已提交
323 324 325
  int8_t  intervalUnit;
  int8_t  slidingUnit;
  int8_t  timezone;
S
sma  
Shengliang Guan 已提交
326
  int32_t dstVgId;  // for stream
S
Shengliang Guan 已提交
327 328 329
  int64_t interval;
  int64_t offset;
  int64_t sliding;
S
sma  
Shengliang Guan 已提交
330
  int32_t exprLen;  // strlen + 1
S
Shengliang Guan 已提交
331
  int32_t tagsFilterLen;
S
sma  
Shengliang Guan 已提交
332 333
  int32_t sqlLen;
  int32_t astLen;
S
Shengliang Guan 已提交
334 335
  char*   expr;
  char*   tagsFilter;
S
sma  
Shengliang Guan 已提交
336 337
  char*   sql;
  char*   ast;
S
Shengliang Guan 已提交
338 339
} SSmaObj;

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;
L
Liu Jicong 已提交
345 346
  int64_t  uid;
  int64_t  dbUid;
S
Shengliang Guan 已提交
347
  int32_t  version;
S
Shengliang Guan 已提交
348
  int32_t  nextColId;
S
sma  
Shengliang Guan 已提交
349 350 351
  float    xFilesFactor;
  int32_t  aggregationMethod;
  int32_t  delay;
S
sma  
Shengliang Guan 已提交
352
  int32_t  ttl;
S
Shengliang Guan 已提交
353 354
  int32_t  numOfColumns;
  int32_t  numOfTags;
S
sma  
Shengliang Guan 已提交
355
  int32_t  numOfSmas;
S
sma  
Shengliang Guan 已提交
356
  int32_t  commentLen;
S
Shengliang Guan 已提交
357
  SSchema* pColumns;
S
Shengliang Guan 已提交
358
  SSchema* pTags;
S
sma  
Shengliang Guan 已提交
359
  SSchema* pSmas;
S
sma  
Shengliang Guan 已提交
360
  char*    comment;
S
Shengliang Guan 已提交
361
  SRWLatch lock;
S
Shengliang Guan 已提交
362
} SStbObj;
S
Shengliang Guan 已提交
363

S
Shengliang Guan 已提交
364
typedef struct {
S
Shengliang Guan 已提交
365 366
  char    name[TSDB_FUNC_NAME_LEN];
  int64_t createdTime;
S
Shengliang Guan 已提交
367 368 369 370 371 372
  int8_t  funcType;
  int8_t  scriptType;
  int8_t  align;
  int8_t  outputType;
  int32_t outputLen;
  int32_t bufSize;
S
Shengliang 已提交
373
  int64_t signature;
S
Shengliang Guan 已提交
374 375
  int32_t commentSize;
  int32_t codeSize;
L
Liu Jicong 已提交
376 377
  char*   pComment;
  char*   pCode;
S
Shengliang Guan 已提交
378
  char    pData[];
S
Shengliang Guan 已提交
379 380
} SFuncObj;

S
Shengliang Guan 已提交
381
typedef struct {
382
  int64_t id;
S
Shengliang Guan 已提交
383 384 385 386 387 388 389
  int8_t  type;
  int8_t  replica;
  int16_t numOfColumns;
  int32_t rowSize;
  int32_t numOfRows;
  int32_t numOfReads;
  int32_t payloadLen;
L
Liu Jicong 已提交
390 391
  void*   pIter;
  SMnode* pMnode;
392
  char    db[TSDB_DB_FNAME_LEN];
S
Shengliang Guan 已提交
393 394 395
  int16_t offset[TSDB_MAX_COLUMNS];
  int32_t bytes[TSDB_MAX_COLUMNS];
  char    payload[];
S
Shengliang Guan 已提交
396 397
} SShowObj;

398 399 400 401 402 403 404 405 406 407 408 409 410 411 412 413 414
typedef struct {
  int64_t id;
  int8_t  type;
  int8_t  replica;
  int16_t numOfColumns;
  int32_t rowSize;
  int32_t numOfRows;
  int32_t numOfReads;
  int32_t payloadLen;
  void*   pIter;
  SMnode* pMnode;
  char    db[TSDB_DB_FNAME_LEN];
  int16_t offset[TSDB_MAX_COLUMNS];
  int32_t bytes[TSDB_MAX_COLUMNS];
  char    payload[];
} SSysTableRetrieveObj;

L
Liu Jicong 已提交
415 416 417
typedef struct {
  int32_t vgId;  // -1 for unassigned
  int32_t status;
L
Liu Jicong 已提交
418
  int32_t epoch;
L
Liu Jicong 已提交
419 420 421 422
  SEpSet  epSet;
  int64_t oldConsumerId;
  int64_t consumerId;  // -1 for unassigned
  char*   qmsg;
L
Liu Jicong 已提交
423 424
} SMqConsumerEp;

L
Liu Jicong 已提交
425
static FORCE_INLINE int32_t tEncodeSMqConsumerEp(void** buf, const SMqConsumerEp* pConsumerEp) {
L
Liu Jicong 已提交
426 427
  int32_t tlen = 0;
  tlen += taosEncodeFixedI32(buf, pConsumerEp->vgId);
L
Liu Jicong 已提交
428
  tlen += taosEncodeFixedI32(buf, pConsumerEp->status);
L
Liu Jicong 已提交
429
  tlen += taosEncodeFixedI32(buf, pConsumerEp->epoch);
430
  tlen += taosEncodeSEpSet(buf, &pConsumerEp->epSet);
L
Liu Jicong 已提交
431
  tlen += taosEncodeFixedI64(buf, pConsumerEp->oldConsumerId);
L
Liu Jicong 已提交
432
  tlen += taosEncodeFixedI64(buf, pConsumerEp->consumerId);
L
Liu Jicong 已提交
433
  tlen += taosEncodeString(buf, pConsumerEp->qmsg);
L
Liu Jicong 已提交
434 435 436 437 438
  return tlen;
}

static FORCE_INLINE void* tDecodeSMqConsumerEp(void** buf, SMqConsumerEp* pConsumerEp) {
  buf = taosDecodeFixedI32(buf, &pConsumerEp->vgId);
L
Liu Jicong 已提交
439
  buf = taosDecodeFixedI32(buf, &pConsumerEp->status);
L
Liu Jicong 已提交
440
  buf = taosDecodeFixedI32(buf, &pConsumerEp->epoch);
441
  buf = taosDecodeSEpSet(buf, &pConsumerEp->epSet);
L
Liu Jicong 已提交
442
  buf = taosDecodeFixedI64(buf, &pConsumerEp->oldConsumerId);
L
Liu Jicong 已提交
443
  buf = taosDecodeFixedI64(buf, &pConsumerEp->consumerId);
L
Liu Jicong 已提交
444
  buf = taosDecodeString(buf, &pConsumerEp->qmsg);
L
Liu Jicong 已提交
445 446 447
  return buf;
}

L
Liu Jicong 已提交
448 449
static FORCE_INLINE void tDeleteSMqConsumerEp(SMqConsumerEp* pConsumerEp) {
  if (pConsumerEp) {
wafwerar's avatar
wafwerar 已提交
450
    taosMemoryFreeClear(pConsumerEp->qmsg);
L
Liu Jicong 已提交
451 452 453
  }
}

L
Liu Jicong 已提交
454 455 456 457
typedef struct {
  int64_t consumerId;
  SArray* vgInfo;  // SArray<SMqConsumerEp>
} SMqSubConsumer;
L
Liu Jicong 已提交
458

L
Liu Jicong 已提交
459 460 461 462
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 已提交
463
  tlen += taosEncodeFixedI32(buf, sz);
L
Liu Jicong 已提交
464 465 466
  for (int32_t i = 0; i < sz; i++) {
    SMqConsumerEp* pCEp = taosArrayGet(pConsumer->vgInfo, i);
    tlen += tEncodeSMqConsumerEp(buf, pCEp);
L
Liu Jicong 已提交
467
  }
L
Liu Jicong 已提交
468 469
  return tlen;
}
L
Liu Jicong 已提交
470

L
Liu Jicong 已提交
471 472 473 474 475 476 477 478 479
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 已提交
480
  }
L
Liu Jicong 已提交
481 482 483 484 485 486 487
  return buf;
}

static FORCE_INLINE void tDeleteSMqSubConsumer(SMqSubConsumer* pSubConsumer) {
  if (pSubConsumer->vgInfo) {
    taosArrayDestroyEx(pSubConsumer->vgInfo, (void (*)(void*))tDeleteSMqConsumerEp);
    pSubConsumer->vgInfo = NULL;
L
Liu Jicong 已提交
488
  }
L
Liu Jicong 已提交
489 490
}

L
Liu Jicong 已提交
491 492 493 494 495 496 497 498 499 500 501 502 503 504 505 506 507 508
typedef struct {
  char    key[TSDB_PARTITION_KEY_LEN];
  int64_t offset;
} SMqOffsetObj;

static FORCE_INLINE int32_t tEncodeSMqOffsetObj(void** buf, const SMqOffsetObj* pOffset) {
  int32_t tlen = 0;
  tlen += taosEncodeString(buf, pOffset->key);
  tlen += taosEncodeFixedI64(buf, pOffset->offset);
  return tlen;
}

static FORCE_INLINE void* tDecodeSMqOffsetObj(void* buf, SMqOffsetObj* pOffset) {
  buf = taosDecodeStringTo(buf, pOffset->key);
  buf = taosDecodeFixedI64(buf, &pOffset->offset);
  return buf;
}

L
Liu Jicong 已提交
509 510 511 512
typedef struct {
  char    key[TSDB_SUBSCRIBE_KEY_LEN];
  int32_t status;
  int32_t vgNum;
L
Liu Jicong 已提交
513 514 515
  SArray* consumers;      // SArray<SMqSubConsumer>
  SArray* lostConsumers;  // SArray<SMqSubConsumer>
  SArray* unassignedVg;   // SArray<SMqConsumerEp>
L
Liu Jicong 已提交
516 517 518
} SMqSubscribeObj;

static FORCE_INLINE SMqSubscribeObj* tNewSubscribeObj() {
wafwerar's avatar
wafwerar 已提交
519
  SMqSubscribeObj* pSub = taosMemoryCalloc(1, sizeof(SMqSubscribeObj));
L
Liu Jicong 已提交
520
  if (pSub == NULL) {
L
Liu Jicong 已提交
521 522
    return NULL;
  }
L
Liu Jicong 已提交
523

L
Liu Jicong 已提交
524 525 526
  pSub->consumers = taosArrayInit(0, sizeof(SMqSubConsumer));
  if (pSub->consumers == NULL) {
    goto _err;
L
Liu Jicong 已提交
527
  }
L
Liu Jicong 已提交
528 529 530 531 532 533

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

L
Liu Jicong 已提交
534
  pSub->unassignedVg = taosArrayInit(0, sizeof(SMqConsumerEp));
L
Liu Jicong 已提交
535
  if (pSub->unassignedVg == NULL) {
L
Liu Jicong 已提交
536
    goto _err;
L
Liu Jicong 已提交
537
  }
L
Liu Jicong 已提交
538 539 540 541 542

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

L
Liu Jicong 已提交
543
  return pSub;
L
Liu Jicong 已提交
544 545

_err:
wafwerar's avatar
wafwerar 已提交
546 547 548 549
  taosMemoryFreeClear(pSub->consumers);
  taosMemoryFreeClear(pSub->lostConsumers);
  taosMemoryFreeClear(pSub->unassignedVg);
  taosMemoryFreeClear(pSub);
L
Liu Jicong 已提交
550
  return NULL;
L
Liu Jicong 已提交
551 552 553 554 555
}

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

L
Liu Jicong 已提交
560
  sz = taosArrayGetSize(pSub->consumers);
L
Liu Jicong 已提交
561 562
  tlen += taosEncodeFixedI32(buf, sz);
  for (int32_t i = 0; i < sz; i++) {
L
Liu Jicong 已提交
563 564
    SMqSubConsumer* pSubConsumer = taosArrayGet(pSub->consumers, i);
    tlen += tEncodeSMqSubConsumer(buf, pSubConsumer);
L
Liu Jicong 已提交
565 566
  }

L
Liu Jicong 已提交
567 568 569 570 571 572 573
  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 已提交
574 575 576 577 578 579 580 581 582 583 584 585
  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 已提交
586 587
  buf = taosDecodeFixedI32(buf, &pSub->vgNum);
  buf = taosDecodeFixedI32(buf, &pSub->status);
L
Liu Jicong 已提交
588 589 590 591

  int32_t sz;

  buf = taosDecodeFixedI32(buf, &sz);
L
Liu Jicong 已提交
592 593
  pSub->consumers = taosArrayInit(sz, sizeof(SMqSubConsumer));
  if (pSub->consumers == NULL) {
L
Liu Jicong 已提交
594 595 596
    return NULL;
  }
  for (int32_t i = 0; i < sz; i++) {
L
Liu Jicong 已提交
597 598 599
    SMqSubConsumer subConsumer = {0};
    buf = tDecodeSMqSubConsumer(buf, &subConsumer);
    taosArrayPush(pSub->consumers, &subConsumer);
L
Liu Jicong 已提交
600 601
  }

L
Liu Jicong 已提交
602 603 604 605 606 607 608 609 610 611 612
  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 已提交
613 614 615 616 617 618
  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 已提交
619 620 621
    SMqConsumerEp consumerEp = {0};
    buf = tDecodeSMqConsumerEp(buf, &consumerEp);
    taosArrayPush(pSub->unassignedVg, &consumerEp);
L
Liu Jicong 已提交
622 623 624
  }
  return buf;
}
L
Liu Jicong 已提交
625

L
Liu Jicong 已提交
626
static FORCE_INLINE void tDeleteSMqSubscribeObj(SMqSubscribeObj* pSub) {
L
Liu Jicong 已提交
627
  if (pSub->consumers) {
L
temp  
Liu Jicong 已提交
628
    //taosArrayDestroyEx(pSub->consumers, (void (*)(void*))tDeleteSMqSubConsumer);
L
Liu Jicong 已提交
629
    // taosArrayDestroy(pSub->consumers);
L
Liu Jicong 已提交
630
    pSub->consumers = NULL;
L
Liu Jicong 已提交
631
  }
L
Liu Jicong 已提交
632

L
Liu Jicong 已提交
633
  if (pSub->unassignedVg) {
L
temp  
Liu Jicong 已提交
634
    //taosArrayDestroyEx(pSub->unassignedVg, (void (*)(void*))tDeleteSMqConsumerEp);
L
Liu Jicong 已提交
635
    // taosArrayDestroy(pSub->unassignedVg);
L
Liu Jicong 已提交
636 637
    pSub->unassignedVg = NULL;
  }
L
Liu Jicong 已提交
638 639
}

L
Liu Jicong 已提交
640
typedef struct {
L
Liu Jicong 已提交
641 642 643 644 645 646 647 648 649 650 651 652 653
  char           name[TSDB_TOPIC_FNAME_LEN];
  char           db[TSDB_DB_FNAME_LEN];
  int64_t        createTime;
  int64_t        updateTime;
  int64_t        uid;
  int64_t        dbUid;
  int32_t        version;
  SRWLatch       lock;
  int32_t        sqlLen;
  char*          sql;
  char*          logicalPlan;
  char*          physicalPlan;
  SSchemaWrapper schema;
L
Liu Jicong 已提交
654 655
} SMqTopicObj;

L
Liu Jicong 已提交
656
typedef struct {
L
Liu Jicong 已提交
657
  int64_t  consumerId;
L
Liu Jicong 已提交
658
  int64_t  connId;
L
Liu Jicong 已提交
659
  SRWLatch lock;
L
Liu Jicong 已提交
660
  char     cgroup[TSDB_CGROUP_LEN];
L
Liu Jicong 已提交
661 662
  SArray*  currentTopics;        // SArray<char*>
  SArray*  recentRemovedTopics;  // SArray<char*>
L
Liu Jicong 已提交
663
  int32_t  epoch;
L
Liu Jicong 已提交
664
  // stat
L
Liu Jicong 已提交
665 666 667 668 669 670 671
  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 已提交
672 673
} SMqConsumerObj;

L
Liu Jicong 已提交
674
static FORCE_INLINE int32_t tEncodeSMqConsumerObj(void** buf, const SMqConsumerObj* pConsumer) {
L
Liu Jicong 已提交
675
  int32_t sz;
L
Liu Jicong 已提交
676 677
  int32_t tlen = 0;
  tlen += taosEncodeFixedI64(buf, pConsumer->consumerId);
L
Liu Jicong 已提交
678
  tlen += taosEncodeFixedI64(buf, pConsumer->connId);
L
Liu Jicong 已提交
679
  tlen += taosEncodeFixedI32(buf, pConsumer->epoch);
L
Liu Jicong 已提交
680
  tlen += taosEncodeFixedI64(buf, pConsumer->pollCnt);
L
Liu Jicong 已提交
681
  tlen += taosEncodeFixedI32(buf, pConsumer->status);
L
Liu Jicong 已提交
682
  tlen += taosEncodeString(buf, pConsumer->cgroup);
L
Liu Jicong 已提交
683 684 685 686 687 688 689 690 691

  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 已提交
692 693
  tlen += taosEncodeFixedI32(buf, sz);
  for (int32_t i = 0; i < sz; i++) {
L
Liu Jicong 已提交
694
    char* topic = taosArrayGetP(pConsumer->recentRemovedTopics, i);
L
Liu Jicong 已提交
695
    tlen += taosEncodeString(buf, topic);
L
Liu Jicong 已提交
696 697 698 699 700
  }
  return tlen;
}

static FORCE_INLINE void* tDecodeSMqConsumerObj(void* buf, SMqConsumerObj* pConsumer) {
L
Liu Jicong 已提交
701
  int32_t sz;
L
Liu Jicong 已提交
702
  buf = taosDecodeFixedI64(buf, &pConsumer->consumerId);
L
Liu Jicong 已提交
703
  buf = taosDecodeFixedI64(buf, &pConsumer->connId);
L
Liu Jicong 已提交
704
  buf = taosDecodeFixedI32(buf, &pConsumer->epoch);
L
Liu Jicong 已提交
705
  buf = taosDecodeFixedI64(buf, &pConsumer->pollCnt);
L
Liu Jicong 已提交
706
  buf = taosDecodeFixedI32(buf, &pConsumer->status);
L
Liu Jicong 已提交
707
  buf = taosDecodeStringTo(buf, pConsumer->cgroup);
L
Liu Jicong 已提交
708 709 710 711 712 713 714 715 716

  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 已提交
717
  buf = taosDecodeFixedI32(buf, &sz);
L
Liu Jicong 已提交
718
  pConsumer->recentRemovedTopics = taosArrayInit(sz, sizeof(SMqConsumerObj));
L
Liu Jicong 已提交
719
  for (int32_t i = 0; i < sz; i++) {
L
Liu Jicong 已提交
720 721
    char* topic;
    buf = taosDecodeString(buf, &topic);
L
Liu Jicong 已提交
722
    taosArrayPush(pConsumer->recentRemovedTopics, &topic);
L
Liu Jicong 已提交
723 724 725 726
  }
  return buf;
}

L
Liu Jicong 已提交
727
typedef struct {
L
Liu Jicong 已提交
728 729
  char     name[TSDB_TOPIC_FNAME_LEN];
  char     db[TSDB_DB_FNAME_LEN];
L
Liu Jicong 已提交
730
  char     outputSTbName[TSDB_TABLE_FNAME_LEN];
L
Liu Jicong 已提交
731 732 733 734 735
  int64_t  createTime;
  int64_t  updateTime;
  int64_t  uid;
  int64_t  dbUid;
  int32_t  version;
L
Liu Jicong 已提交
736
  int32_t  vgNum;
L
Liu Jicong 已提交
737 738 739
  SRWLatch lock;
  int8_t   status;
  // int32_t  sqlLen;
L
Liu Jicong 已提交
740 741 742
  int8_t         createdBy;      // STREAM_CREATED_BY__USER or SMA
  int32_t        fixedSinkVgId;  // 0 for shuffle
  int64_t        smaId;          // 0 for unused
L
fix  
Liu Jicong 已提交
743 744 745
  int8_t         trigger;
  int32_t        triggerParam;
  int64_t        waterMark;
L
Liu Jicong 已提交
746 747 748 749 750
  char*          sql;
  char*          logicalPlan;
  char*          physicalPlan;
  SArray*        tasks;  // SArray<SArray<SStreamTask>>
  SSchemaWrapper outputSchema;
L
Liu Jicong 已提交
751 752 753 754
} SStreamObj;

int32_t tEncodeSStreamObj(SCoder* pEncoder, const SStreamObj* pObj);
int32_t tDecodeSStreamObj(SCoder* pDecoder, SStreamObj* pObj);
L
Liu Jicong 已提交
755

S
Shengliang Guan 已提交
756 757 758
#ifdef __cplusplus
}
#endif
S
Shengliang Guan 已提交
759

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