vnodeSvr.c 50.8 KB
Newer Older
H
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/>.
 */

H
Hongze Cheng 已提交
16 17
#include "tencode.h"
#include "tmsg.h"
H
Hongze Cheng 已提交
18
#include "vnd.h"
H
Hongze Cheng 已提交
19 20
#include "vnode.h"
#include "vnodeInt.h"
H
Hongze Cheng 已提交
21

22 23 24 25 26 27 28 29 30
static int32_t vnodeProcessCreateStbReq(SVnode *pVnode, int64_t version, void *pReq, int32_t len, SRpcMsg *pRsp);
static int32_t vnodeProcessAlterStbReq(SVnode *pVnode, int64_t version, void *pReq, int32_t len, SRpcMsg *pRsp);
static int32_t vnodeProcessDropStbReq(SVnode *pVnode, int64_t version, void *pReq, int32_t len, SRpcMsg *pRsp);
static int32_t vnodeProcessCreateTbReq(SVnode *pVnode, int64_t version, void *pReq, int32_t len, SRpcMsg *pRsp);
static int32_t vnodeProcessAlterTbReq(SVnode *pVnode, int64_t version, void *pReq, int32_t len, SRpcMsg *pRsp);
static int32_t vnodeProcessDropTbReq(SVnode *pVnode, int64_t version, void *pReq, int32_t len, SRpcMsg *pRsp);
static int32_t vnodeProcessSubmitReq(SVnode *pVnode, int64_t version, void *pReq, int32_t len, SRpcMsg *pRsp);
static int32_t vnodeProcessCreateTSmaReq(SVnode *pVnode, int64_t version, void *pReq, int32_t len, SRpcMsg *pRsp);
static int32_t vnodeProcessAlterConfirmReq(SVnode *pVnode, int64_t version, void *pReq, int32_t len, SRpcMsg *pRsp);
31
static int32_t vnodeProcessAlterConfigReq(SVnode *pVnode, int64_t version, void *pReq, int32_t len, SRpcMsg *pRsp);
32
static int32_t vnodeProcessDropTtlTbReq(SVnode *pVnode, int64_t version, void *pReq, int32_t len, SRpcMsg *pRsp);
S
Shengliang Guan 已提交
33
static int32_t vnodeProcessTrimReq(SVnode *pVnode, int64_t version, void *pReq, int32_t len, SRpcMsg *pRsp);
H
Hongze Cheng 已提交
34
static int32_t vnodeProcessDeleteReq(SVnode *pVnode, int64_t version, void *pReq, int32_t len, SRpcMsg *pRsp);
35
static int32_t vnodeProcessBatchDeleteReq(SVnode *pVnode, int64_t version, void *pReq, int32_t len, SRpcMsg *pRsp);
dengyihao's avatar
dengyihao 已提交
36 37
static int32_t vnodeProcessCreateIndexReq(SVnode *pVnode, int64_t version, void *pReq, int32_t len, SRpcMsg *pRsp);
static int32_t vnodeProcessDropIndexReq(SVnode *pVnode, int64_t version, void *pReq, int32_t len, SRpcMsg *pRsp);
38
static int32_t vnodeProcessCompactVnodeReq(SVnode *pVnode, int64_t version, void *pReq, int32_t len, SRpcMsg *pRsp);
H
Hongze Cheng 已提交
39

H
Hongze Cheng 已提交
40
static int32_t vnodePreprocessCreateTableReq(SVnode *pVnode, SDecoder *pCoder, int64_t ctime, int64_t *pUid) {
H
Hongze Cheng 已提交
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 76 77 78
  int32_t code = 0;
  int32_t lino = 0;

  if (tStartDecode(pCoder) < 0) {
    code = TSDB_CODE_INVALID_MSG;
    TSDB_CHECK_CODE(code, lino, _exit);
  }

  // flags
  if (tDecodeI32v(pCoder, NULL) < 0) {
    code = TSDB_CODE_INVALID_MSG;
    TSDB_CHECK_CODE(code, lino, _exit);
  }

  // name
  char *name = NULL;
  if (tDecodeCStr(pCoder, &name) < 0) {
    code = TSDB_CODE_INVALID_MSG;
    TSDB_CHECK_CODE(code, lino, _exit);
  }

  // uid
  int64_t uid = metaGetTableEntryUidByName(pVnode->pMeta, name);
  if (uid == 0) {
    uid = tGenIdPI64();
  }
  *(int64_t *)(pCoder->data + pCoder->pos) = uid;

  // ctime
  *(int64_t *)(pCoder->data + pCoder->pos + 8) = ctime;

  tEndDecode(pCoder);

_exit:
  if (code) {
    vError("vgId:%d %s failed at line %d since %s", TD_VID(pVnode), __func__, lino, tstrerror(code));
  } else {
    vTrace("vgId:%d %s done, table:%s uid generated:%" PRId64, TD_VID(pVnode), __func__, name, uid);
H
Hongze Cheng 已提交
79
    if (pUid) *pUid = uid;
H
Hongze Cheng 已提交
80 81 82
  }
  return code;
}
H
Hongze Cheng 已提交
83 84 85 86 87
static int32_t vnodePreProcessCreateTableMsg(SVnode *pVnode, SRpcMsg *pMsg) {
  int32_t code = 0;
  int32_t lino = 0;

  int64_t  ctime = taosGetTimestampMs();
H
Hongze Cheng 已提交
88
  SDecoder dc = {0};
H
Hongze Cheng 已提交
89
  int32_t  nReqs;
H
Hongze Cheng 已提交
90

H
Hongze Cheng 已提交
91 92 93 94 95 96 97 98 99 100 101
  tDecoderInit(&dc, (uint8_t *)pMsg->pCont + sizeof(SMsgHead), pMsg->contLen - sizeof(SMsgHead));
  if (tStartDecode(&dc) < 0) {
    code = TSDB_CODE_INVALID_MSG;
    return code;
  }

  if (tDecodeI32v(&dc, &nReqs) < 0) {
    code = TSDB_CODE_INVALID_MSG;
    TSDB_CHECK_CODE(code, lino, _exit);
  }
  for (int32_t iReq = 0; iReq < nReqs; iReq++) {
H
Hongze Cheng 已提交
102
    code = vnodePreprocessCreateTableReq(pVnode, &dc, ctime, NULL);
H
Hongze Cheng 已提交
103
    TSDB_CHECK_CODE(code, lino, _exit);
H
Hongze Cheng 已提交
104 105 106 107 108 109 110 111
  }

  tEndDecode(&dc);

_exit:
  tDecoderClear(&dc);
  return code;
}
H
Hongze Cheng 已提交
112 113
extern int64_t tsMaxKeyByPrecision[];
static int32_t vnodePreProcessSubmitTbData(SVnode *pVnode, SDecoder *pCoder, int64_t ctime) {
H
Hongze Cheng 已提交
114 115 116
  int32_t code = 0;
  int32_t lino = 0;

H
Hongze Cheng 已提交
117 118 119 120
  if (tStartDecode(pCoder) < 0) {
    code = TSDB_CODE_INVALID_MSG;
    TSDB_CHECK_CODE(code, lino, _exit);
  }
H
Hongze Cheng 已提交
121

H
Hongze Cheng 已提交
122 123 124 125 126
  SSubmitTbData submitTbData;
  if (tDecodeI32v(pCoder, &submitTbData.flags) < 0) {
    code = TSDB_CODE_INVALID_MSG;
    TSDB_CHECK_CODE(code, lino, _exit);
  }
H
Hongze Cheng 已提交
127

H
Hongze Cheng 已提交
128
  int64_t uid;
H
Hongze Cheng 已提交
129
  if (submitTbData.flags & SUBMIT_REQ_AUTO_CREATE_TABLE) {
H
Hongze Cheng 已提交
130
    code = vnodePreprocessCreateTableReq(pVnode, pCoder, ctime, &uid);
H
Hongze Cheng 已提交
131 132 133 134 135 136 137 138
    TSDB_CHECK_CODE(code, lino, _exit);
  }

  // submit data
  if (tDecodeI64(pCoder, &submitTbData.suid) < 0) {
    code = TSDB_CODE_INVALID_MSG;
    TSDB_CHECK_CODE(code, lino, _exit);
  }
H
Hongze Cheng 已提交
139 140 141

  if (submitTbData.flags & SUBMIT_REQ_AUTO_CREATE_TABLE) {
    *(int64_t *)(pCoder->data + pCoder->pos) = uid;
H
Hongze Cheng 已提交
142
    pCoder->pos += sizeof(int64_t);
H
Hongze Cheng 已提交
143
  } else {
H
Hongze Cheng 已提交
144 145 146 147
    if (tDecodeI64(pCoder, &submitTbData.uid) < 0) {
      code = TSDB_CODE_INVALID_MSG;
      TSDB_CHECK_CODE(code, lino, _exit);
    }
H
Hongze Cheng 已提交
148
  }
H
Hongze Cheng 已提交
149

H
Hongze Cheng 已提交
150
  if (tDecodeI32v(pCoder, &submitTbData.sver) < 0) {
H
Hongze Cheng 已提交
151 152 153 154
    code = TSDB_CODE_INVALID_MSG;
    TSDB_CHECK_CODE(code, lino, _exit);
  }

H
Hongze Cheng 已提交
155 156 157 158 159 160 161 162 163 164 165 166
  // scan and check
  TSKEY now = ctime;
  if (pVnode->config.tsdbCfg.precision == TSDB_TIME_PRECISION_MICRO) {
    now *= 1000;
  } else if (pVnode->config.tsdbCfg.precision == TSDB_TIME_PRECISION_NANO) {
    now *= 1000000;
  }
  TSKEY minKey = now - tsTickPerMin[pVnode->config.tsdbCfg.precision] * pVnode->config.tsdbCfg.keep2;
  TSKEY maxKey = tsMaxKeyByPrecision[pVnode->config.tsdbCfg.precision];
  if (submitTbData.flags & SUBMIT_REQ_COLUMN_DATA_FORMAT) {
    uint64_t nColData;
    if (tDecodeU64v(pCoder, &nColData) < 0) {
H
Hongze Cheng 已提交
167
      code = TSDB_CODE_INVALID_MSG;
H
Hongze Cheng 已提交
168
      goto _exit;
H
Hongze Cheng 已提交
169 170
    }

H
Hongze Cheng 已提交
171 172 173
    SColData colData = {0};
    pCoder->pos += tGetColData(pCoder->data + pCoder->pos, &colData);

H
Hongze Cheng 已提交
174 175 176 177 178
    if (colData.flag != HAS_VALUE) {
      code = TSDB_CODE_INVALID_MSG;
      goto _exit;
    }

H
Hongze Cheng 已提交
179 180 181 182 183 184 185 186 187
    for (int32_t iRow = 0; iRow < colData.nVal; iRow++) {
      if (((TSKEY *)colData.pData)[iRow] < minKey || ((TSKEY *)colData.pData)[iRow] > maxKey) {
        code = TSDB_CODE_TDB_TIMESTAMP_OUT_OF_RANGE;
        goto _exit;
      }
    }
  } else {
    uint64_t nRow;
    if (tDecodeU64v(pCoder, &nRow) < 0) {
H
Hongze Cheng 已提交
188
      code = TSDB_CODE_INVALID_MSG;
H
Hongze Cheng 已提交
189
      goto _exit;
H
Hongze Cheng 已提交
190
    }
H
Hongze Cheng 已提交
191

H
Hongze Cheng 已提交
192 193 194
    for (int32_t iRow = 0; iRow < nRow; ++iRow) {
      SRow *pRow = (SRow *)(pCoder->data + pCoder->pos);
      pCoder->pos += pRow->len;
H
Hongze Cheng 已提交
195

H
Hongze Cheng 已提交
196 197 198
      if (pRow->ts < minKey || pRow->ts > maxKey) {
        code = TSDB_CODE_TDB_TIMESTAMP_OUT_OF_RANGE;
        goto _exit;
H
Hongze Cheng 已提交
199
      }
H
Hongze Cheng 已提交
200 201
    }
  }
H
Hongze Cheng 已提交
202

H
Hongze Cheng 已提交
203
  tEndDecode(pCoder);
H
Hongze Cheng 已提交
204

H
Hongze Cheng 已提交
205
_exit:
H
Hongze Cheng 已提交
206
  return code;
H
Hongze Cheng 已提交
207 208 209 210
}
static int32_t vnodePreProcessSubmitMsg(SVnode *pVnode, SRpcMsg *pMsg) {
  int32_t code = 0;
  int32_t lino = 0;
H
Hongze Cheng 已提交
211

H
Hongze Cheng 已提交
212
  SDecoder *pCoder = &(SDecoder){0};
H
Hongze Cheng 已提交
213

X
Xiaoyu Wang 已提交
214
  tDecoderInit(pCoder, (uint8_t *)pMsg->pCont + sizeof(SSubmitReq2Msg), pMsg->contLen - sizeof(SSubmitReq2Msg));
H
Hongze Cheng 已提交
215

H
Hongze Cheng 已提交
216 217 218 219
  if (tStartDecode(pCoder) < 0) {
    code = TSDB_CODE_INVALID_MSG;
    TSDB_CHECK_CODE(code, lino, _exit);
  }
H
Hongze Cheng 已提交
220

H
Hongze Cheng 已提交
221 222 223 224 225
  uint64_t nSubmitTbData;
  if (tDecodeU64v(pCoder, &nSubmitTbData) < 0) {
    code = TSDB_CODE_INVALID_MSG;
    TSDB_CHECK_CODE(code, lino, _exit);
  }
H
Hongze Cheng 已提交
226

H
Hongze Cheng 已提交
227 228 229 230
  int64_t ctime = taosGetTimestampMs();
  for (int32_t i = 0; i < nSubmitTbData; i++) {
    code = vnodePreProcessSubmitTbData(pVnode, pCoder, ctime);
    TSDB_CHECK_CODE(code, lino, _exit);
H
Hongze Cheng 已提交
231
  }
H
Hongze Cheng 已提交
232

H
Hongze Cheng 已提交
233
  tEndDecode(pCoder);
H
Hongze Cheng 已提交
234

H
Hongze Cheng 已提交
235
_exit:
H
Hongze Cheng 已提交
236
  tDecoderClear(pCoder);
H
Hongze Cheng 已提交
237 238
  return code;
}
H
Hongze Cheng 已提交
239

H
Hongze Cheng 已提交
240 241 242 243 244 245 246 247 248 249 250 251 252 253 254 255
static int32_t vnodePreProcessDeleteMsg(SVnode *pVnode, SRpcMsg *pMsg) {
  int32_t code = 0;

  int32_t     size;
  int32_t     ret;
  uint8_t    *pCont;
  SEncoder   *pCoder = &(SEncoder){0};
  SDeleteRes  res = {0};
  SReadHandle handle = {.meta = pVnode->pMeta, .config = &pVnode->config, .vnode = pVnode, .pMsgCb = &pVnode->msgCb};

  code = qWorkerProcessDeleteMsg(&handle, pVnode->pQuery, pMsg, &res);
  if (code) goto _exit;

  // malloc and encode
  tEncodeSize(tEncodeDeleteRes, &res, size, ret);
  pCont = rpcMallocCont(size + sizeof(SMsgHead));
H
Hongze Cheng 已提交
256

H
Hongze Cheng 已提交
257 258
  ((SMsgHead *)pCont)->contLen = size + sizeof(SMsgHead);
  ((SMsgHead *)pCont)->vgId = TD_VID(pVnode);
H
Hongze Cheng 已提交
259

H
Hongze Cheng 已提交
260 261 262
  tEncoderInit(pCoder, pCont + sizeof(SMsgHead), size);
  tEncodeDeleteRes(pCoder, &res);
  tEncoderClear(pCoder);
H
Hongze Cheng 已提交
263

H
Hongze Cheng 已提交
264 265 266
  rpcFreeCont(pMsg->pCont);
  pMsg->pCont = pCont;
  pMsg->contLen = size + sizeof(SMsgHead);
H
Hongze Cheng 已提交
267

H
Hongze Cheng 已提交
268 269 270 271 272
  taosArrayDestroy(res.uidList);

_exit:
  return code;
}
H
Hongze Cheng 已提交
273

H
Hongze Cheng 已提交
274 275 276 277 278 279 280 281 282 283 284 285
int32_t vnodePreProcessWriteMsg(SVnode *pVnode, SRpcMsg *pMsg) {
  int32_t code = 0;

  switch (pMsg->msgType) {
    case TDMT_VND_CREATE_TABLE: {
      code = vnodePreProcessCreateTableMsg(pVnode, pMsg);
    } break;
    case TDMT_VND_SUBMIT: {
      code = vnodePreProcessSubmitMsg(pVnode, pMsg);
    } break;
    case TDMT_VND_DELETE: {
      code = vnodePreProcessDeleteMsg(pVnode, pMsg);
H
Hongze Cheng 已提交
286
    } break;
H
Hongze Cheng 已提交
287 288 289
    default:
      break;
  }
H
Hongze Cheng 已提交
290

H
Hongze Cheng 已提交
291 292 293 294 295
_exit:
  if (code) {
    vError("vgId%d failed to preprocess write request since %s, msg type:%d", TD_VID(pVnode), tstrerror(code),
           pMsg->msgType);
  }
H
Hongze Cheng 已提交
296
  return code;
H
Hongze Cheng 已提交
297 298
}

299
int32_t vnodeProcessWriteMsg(SVnode *pVnode, SRpcMsg *pMsg, int64_t version, SRpcMsg *pRsp) {
300 301 302 303
  void   *ptr = NULL;
  void   *pReq;
  int32_t len;
  int32_t ret;
304
  /*
305
  if (!pVnode->inUse) {
306
    terrno = TSDB_CODE_VND_NO_AVAIL_BUFPOOL;
H
Hongze Cheng 已提交
307 308
    vError("vgId:%d, not ready to write since %s", TD_VID(pVnode), terrstr());
    return -1;
309
  }
310
  */
311 312 313
  if (version <= pVnode->state.applied) {
    vError("vgId:%d, duplicate write request. version: %" PRId64 ", applied: %" PRId64 "", TD_VID(pVnode), version,
           pVnode->state.applied);
314
    terrno = TSDB_CODE_VND_DUP_REQUEST;
315 316 317
    return -1;
  }

318
  vDebug("vgId:%d, start to process write request %s, index:%" PRId64, TD_VID(pVnode), TMSG_INFO(pMsg->msgType),
H
Hongze Cheng 已提交
319
         version);
H
Hongze Cheng 已提交
320

321 322
  ASSERT(pVnode->state.applyTerm <= pMsg->info.conn.applyTerm);
  ASSERT(pVnode->state.applied + 1 == version);
323

H
Hongze Cheng 已提交
324
  pVnode->state.applied = version;
H
Hongze Cheng 已提交
325
  pVnode->state.applyTerm = pMsg->info.conn.applyTerm;
H
Hongze Cheng 已提交
326

327 328
  if (!syncUtilUserCommit(pMsg->msgType)) goto _exit;

L
Liu Jicong 已提交
329
  if (pMsg->msgType == TDMT_VND_STREAM_RECOVER_BLOCKING_STAGE || pMsg->msgType == TDMT_STREAM_TASK_CHECK_RSP) {
L
Liu Jicong 已提交
330 331 332
    if (tqCheckLogInWal(pVnode->pTq, version)) return 0;
  }

H
Hongze Cheng 已提交
333 334 335
  // skip header
  pReq = POINTER_SHIFT(pMsg->pCont, sizeof(SMsgHead));
  len = pMsg->contLen - sizeof(SMsgHead);
B
Benguang Zhao 已提交
336
  bool needCommit = false;
H
Hongze Cheng 已提交
337 338

  switch (pMsg->msgType) {
H
Hongze Cheng 已提交
339
    /* META */
H
Hongze Cheng 已提交
340
    case TDMT_VND_CREATE_STB:
H
Hongze Cheng 已提交
341
      if (vnodeProcessCreateStbReq(pVnode, version, pReq, len, pRsp) < 0) goto _err;
H
Hongze Cheng 已提交
342
      break;
H
Hongze Cheng 已提交
343
    case TDMT_VND_ALTER_STB:
H
Hongze Cheng 已提交
344
      if (vnodeProcessAlterStbReq(pVnode, version, pReq, len, pRsp) < 0) goto _err;
H
Hongze Cheng 已提交
345
      break;
H
Hongze Cheng 已提交
346
    case TDMT_VND_DROP_STB:
H
Hongze Cheng 已提交
347
      if (vnodeProcessDropStbReq(pVnode, version, pReq, len, pRsp) < 0) goto _err;
H
Hongze Cheng 已提交
348
      break;
H
Hongze Cheng 已提交
349
    case TDMT_VND_CREATE_TABLE:
H
Hongze Cheng 已提交
350
      if (vnodeProcessCreateTbReq(pVnode, version, pReq, len, pRsp) < 0) goto _err;
H
Hongze Cheng 已提交
351 352
      break;
    case TDMT_VND_ALTER_TABLE:
H
Hongze Cheng 已提交
353
      if (vnodeProcessAlterTbReq(pVnode, version, pReq, len, pRsp) < 0) goto _err;
H
Hongze Cheng 已提交
354
      break;
H
Hongze Cheng 已提交
355
    case TDMT_VND_DROP_TABLE:
H
Hongze Cheng 已提交
356
      if (vnodeProcessDropTbReq(pVnode, version, pReq, len, pRsp) < 0) goto _err;
H
Hongze Cheng 已提交
357
      break;
358
    case TDMT_VND_DROP_TTL_TABLE:
359
      if (vnodeProcessDropTtlTbReq(pVnode, version, pReq, len, pRsp) < 0) goto _err;
360
      break;
S
Shengliang Guan 已提交
361 362 363 364
    case TDMT_VND_TRIM:
      if (vnodeProcessTrimReq(pVnode, version, pReq, len, pRsp) < 0) goto _err;
      break;
    case TDMT_VND_CREATE_SMA:
365
      if (vnodeProcessCreateTSmaReq(pVnode, version, pReq, len, pRsp) < 0) goto _err;
S
Shengliang Guan 已提交
366
      break;
H
Hongze Cheng 已提交
367
    /* TSDB */
H
Hongze Cheng 已提交
368
    case TDMT_VND_SUBMIT:
369
      if (vnodeProcessSubmitReq(pVnode, version, pMsg->pCont, pMsg->contLen, pRsp) < 0) goto _err;
H
Hongze Cheng 已提交
370
      break;
D
dapan1121 已提交
371
    case TDMT_VND_DELETE:
H
Hongze Cheng 已提交
372
      if (vnodeProcessDeleteReq(pVnode, version, pReq, len, pRsp) < 0) goto _err;
D
dapan1121 已提交
373
      break;
374 375 376
    case TDMT_VND_BATCH_DEL:
      if (vnodeProcessBatchDeleteReq(pVnode, version, pReq, len, pRsp) < 0) goto _err;
      break;
H
Hongze Cheng 已提交
377
    /* TQ */
L
Liu Jicong 已提交
378
    case TDMT_VND_TMQ_SUBSCRIBE:
379
      if (tqProcessSubscribeReq(pVnode->pTq, version, pReq, len) < 0) {
380
        goto _err;
L
Liu Jicong 已提交
381 382
      }
      break;
L
Liu Jicong 已提交
383 384
    case TDMT_VND_TMQ_DELETE_SUB:
      if (tqProcessDeleteSubReq(pVnode->pTq, version, pMsg->pCont, pMsg->contLen) < 0) {
385 386 387
        goto _err;
      }
      break;
L
Liu Jicong 已提交
388
    case TDMT_VND_TMQ_COMMIT_OFFSET:
389
      if (tqProcessOffsetCommitReq(pVnode->pTq, version, pReq, pMsg->contLen - sizeof(SMsgHead)) < 0) {
390
        goto _err;
L
Liu Jicong 已提交
391 392
      }
      break;
L
Liu Jicong 已提交
393
    case TDMT_VND_TMQ_ADD_CHECKINFO:
394
      if (tqProcessAddCheckInfoReq(pVnode->pTq, version, pReq, len) < 0) {
395 396 397
        goto _err;
      }
      break;
L
Liu Jicong 已提交
398
    case TDMT_VND_TMQ_DEL_CHECKINFO:
399
      if (tqProcessDelCheckInfoReq(pVnode->pTq, version, pReq, len) < 0) {
L
Liu Jicong 已提交
400 401 402
        goto _err;
      }
      break;
403
    case TDMT_STREAM_TASK_DEPLOY: {
404
      if (tqProcessTaskDeployReq(pVnode->pTq, version, pReq, len) < 0) {
405
        goto _err;
H
Hongze Cheng 已提交
406 407
      }
    } break;
L
Liu Jicong 已提交
408
    case TDMT_STREAM_TASK_DROP: {
409
      if (tqProcessTaskDropReq(pVnode->pTq, version, pMsg->pCont, pMsg->contLen) < 0) {
L
Liu Jicong 已提交
410 411 412
        goto _err;
      }
    } break;
L
Liu Jicong 已提交
413
    case TDMT_VND_STREAM_RECOVER_BLOCKING_STAGE: {
L
Liu Jicong 已提交
414 415 416 417
      if (tqProcessTaskRecover2Req(pVnode->pTq, version, pMsg->pCont, pMsg->contLen) < 0) {
        goto _err;
      }
    } break;
418 419 420 421 422
    case TDMT_STREAM_TASK_CHECK_RSP: {
      if (tqProcessStreamTaskCheckRsp(pVnode->pTq, version, pReq, len) < 0) {
        goto _err;
      }
    } break;
423 424 425
    case TDMT_VND_ALTER_CONFIRM:
      vnodeProcessAlterConfirmReq(pVnode, version, pReq, len, pRsp);
      break;
426
    case TDMT_VND_ALTER_CONFIG:
427
      vnodeProcessAlterConfigReq(pVnode, version, pReq, len, pRsp);
S
Shengliang Guan 已提交
428
      break;
H
Hongze Cheng 已提交
429
    case TDMT_VND_COMMIT:
B
Benguang Zhao 已提交
430 431
      needCommit = true;
      break;
dengyihao's avatar
dengyihao 已提交
432 433 434 435 436 437
    case TDMT_VND_CREATE_INDEX:
      vnodeProcessCreateIndexReq(pVnode, version, pReq, len, pRsp);
      break;
    case TDMT_VND_DROP_INDEX:
      vnodeProcessDropIndexReq(pVnode, version, pReq, len, pRsp);
      break;
H
Hongze Cheng 已提交
438
    case TDMT_VND_COMPACT:
439
      vnodeProcessCompactVnodeReq(pVnode, version, pReq, len, pRsp);
H
Hongze Cheng 已提交
440
      goto _exit;
H
Hongze Cheng 已提交
441
    default:
L
Liu Jicong 已提交
442 443
      vError("vgId:%d, unprocessed msg, %d", TD_VID(pVnode), pMsg->msgType);
      return -1;
H
Hongze Cheng 已提交
444 445
  }

S
Shengliang Guan 已提交
446 447
  vTrace("vgId:%d, process %s request, code:0x%x index:%" PRId64, TD_VID(pVnode), TMSG_INFO(pMsg->msgType), pRsp->code,
         version);
H
Hongze Cheng 已提交
448

449 450
  walApplyVer(pVnode->pWal, version);

L
fix log  
Liu Jicong 已提交
451
  /*vInfo("vgId:%d, push msg begin", pVnode->config.vgId);*/
452
  if (tqPushMsg(pVnode->pTq, pMsg->pCont, pMsg->contLen, pMsg->msgType, version) < 0) {
L
fix log  
Liu Jicong 已提交
453
    /*vInfo("vgId:%d, push msg end", pVnode->config.vgId);*/
S
Shengliang Guan 已提交
454
    vError("vgId:%d, failed to push msg to TQ since %s", TD_VID(pVnode), tstrerror(terrno));
455 456
    return -1;
  }
L
fix log  
Liu Jicong 已提交
457
  /*vInfo("vgId:%d, push msg end", pVnode->config.vgId);*/
458

H
Hongze Cheng 已提交
459
  // commit if need
B
Benguang Zhao 已提交
460
  if (needCommit) {
S
Shengliang Guan 已提交
461
    vInfo("vgId:%d, commit at version %" PRId64, TD_VID(pVnode), version);
462 463 464 465
    if (vnodeAsyncCommit(pVnode) < 0) {
      vError("vgId:%d, failed to vnode async commit since %s.", TD_VID(pVnode), tstrerror(terrno));
      goto _err;
    }
H
Hongze Cheng 已提交
466 467

    // start a new one
468
    if (vnodeBegin(pVnode) < 0) {
H
Hongze Cheng 已提交
469 470
      vError("vgId:%d, failed to begin vnode since %s.", TD_VID(pVnode), tstrerror(terrno));
      goto _err;
471
    }
H
Hongze Cheng 已提交
472 473
  }

H
Hongze Cheng 已提交
474
_exit:
H
Hongze Cheng 已提交
475
  return 0;
H
Hongze Cheng 已提交
476 477

_err:
S
Shengliang Guan 已提交
478
  vError("vgId:%d, process %s request failed since %s, version:%" PRId64, TD_VID(pVnode), TMSG_INFO(pMsg->msgType),
H
Hongze Cheng 已提交
479 480
         tstrerror(terrno), version);
  return -1;
H
Hongze Cheng 已提交
481 482
}

483
int32_t vnodePreprocessQueryMsg(SVnode *pVnode, SRpcMsg *pMsg) {
D
dapan1121 已提交
484
  if (TDMT_SCH_QUERY != pMsg->msgType && TDMT_SCH_MERGE_QUERY != pMsg->msgType) {
D
dapan1121 已提交
485 486 487
    return 0;
  }

488
  return qWorkerPreprocessQueryMsg(pVnode->pQuery, pMsg, TDMT_SCH_QUERY == pMsg->msgType);
D
dapan1121 已提交
489 490 491
}

int32_t vnodeProcessQueryMsg(SVnode *pVnode, SRpcMsg *pMsg) {
D
dapan 已提交
492
  vTrace("message in vnode query queue is processing");
493
  if ((pMsg->msgType == TDMT_SCH_QUERY) && !syncIsReadyForRead(pVnode->sync)) {
494
    vnodeRedirectRpcMsg(pVnode, pMsg, terrno);
495 496 497
    return 0;
  }

498
  SReadHandle handle = {.meta = pVnode->pMeta, .config = &pVnode->config, .vnode = pVnode, .pMsgCb = &pVnode->msgCb};
H
Hongze Cheng 已提交
499
  switch (pMsg->msgType) {
D
dapan1121 已提交
500
    case TDMT_SCH_QUERY:
D
dapan1121 已提交
501
    case TDMT_SCH_MERGE_QUERY:
D
dapan1121 已提交
502
      return qWorkerProcessQueryMsg(&handle, pVnode->pQuery, pMsg, 0);
D
dapan1121 已提交
503
    case TDMT_SCH_QUERY_CONTINUE:
D
dapan1121 已提交
504
      return qWorkerProcessCQueryMsg(&handle, pVnode->pQuery, pMsg, 0);
H
Hongze Cheng 已提交
505 506
    default:
      vError("unknown msg type:%d in query queue", pMsg->msgType);
507
      return TSDB_CODE_APP_ERROR;
H
Hongze Cheng 已提交
508 509 510
  }
}

511
int32_t vnodeProcessFetchMsg(SVnode *pVnode, SRpcMsg *pMsg, SQueueInfo *pInfo) {
S
Shengliang Guan 已提交
512
  vTrace("vgId:%d, msg:%p in fetch queue is processing", pVnode->config.vgId, pMsg);
L
Liu Jicong 已提交
513 514
  if ((pMsg->msgType == TDMT_SCH_FETCH || pMsg->msgType == TDMT_VND_TABLE_META || pMsg->msgType == TDMT_VND_TABLE_CFG ||
       pMsg->msgType == TDMT_VND_BATCH_META) &&
515
      !syncIsReadyForRead(pVnode->sync)) {
516
    vnodeRedirectRpcMsg(pVnode, pMsg, terrno);
517 518 519
    return 0;
  }

L
Liu Jicong 已提交
520
  if (pMsg->msgType == TDMT_VND_TMQ_CONSUME && !pVnode->restored) {
521
    vnodeRedirectRpcMsg(pVnode, pMsg, TSDB_CODE_SYN_RESTORING);
L
Liu Jicong 已提交
522 523 524
    return 0;
  }

H
Hongze Cheng 已提交
525
  switch (pMsg->msgType) {
D
dapan1121 已提交
526
    case TDMT_SCH_FETCH:
D
dapan1121 已提交
527
    case TDMT_SCH_MERGE_FETCH:
D
dapan1121 已提交
528
      return qWorkerProcessFetchMsg(pVnode, pVnode->pQuery, pMsg, 0);
D
dapan1121 已提交
529
    case TDMT_SCH_FETCH_RSP:
D
dapan1121 已提交
530
      return qWorkerProcessRspMsg(pVnode, pVnode->pQuery, pMsg, 0);
L
Liu Jicong 已提交
531 532
    // case TDMT_SCH_CANCEL_TASK:
    //   return qWorkerProcessCancelMsg(pVnode, pVnode->pQuery, pMsg, 0);
D
dapan1121 已提交
533
    case TDMT_SCH_DROP_TASK:
D
dapan1121 已提交
534
      return qWorkerProcessDropMsg(pVnode, pVnode->pQuery, pMsg, 0);
D
dapan1121 已提交
535
    case TDMT_SCH_QUERY_HEARTBEAT:
D
dapan1121 已提交
536
      return qWorkerProcessHbMsg(pVnode, pVnode->pQuery, pMsg, 0);
H
Hongze Cheng 已提交
537
    case TDMT_VND_TABLE_META:
D
dapan1121 已提交
538
      return vnodeGetTableMeta(pVnode, pMsg, true);
D
dapan1121 已提交
539
    case TDMT_VND_TABLE_CFG:
D
dapan1121 已提交
540 541 542
      return vnodeGetTableCfg(pVnode, pMsg, true);
    case TDMT_VND_BATCH_META:
      return vnodeGetBatchMeta(pVnode, pMsg);
L
Liu Jicong 已提交
543
    case TDMT_VND_TMQ_CONSUME:
L
Liu Jicong 已提交
544
      return tqProcessPollReq(pVnode->pTq, pMsg);
545 546
    case TDMT_STREAM_TASK_RUN:
      return tqProcessTaskRunReq(pVnode->pTq, pMsg);
547
#if 1
548
    case TDMT_STREAM_TASK_DISPATCH:
L
Liu Jicong 已提交
549
      return tqProcessTaskDispatchReq(pVnode->pTq, pMsg, true);
L
Liu Jicong 已提交
550
#endif
551 552
    case TDMT_STREAM_TASK_CHECK:
      return tqProcessStreamTaskCheckReq(pVnode->pTq, pMsg);
553
    case TDMT_STREAM_TASK_DISPATCH_RSP:
L
Liu Jicong 已提交
554
      return tqProcessTaskDispatchRsp(pVnode->pTq, pMsg);
L
Liu Jicong 已提交
555 556
    case TDMT_STREAM_RETRIEVE:
      return tqProcessTaskRetrieveReq(pVnode->pTq, pMsg);
L
Liu Jicong 已提交
557 558
    case TDMT_STREAM_RETRIEVE_RSP:
      return tqProcessTaskRetrieveRsp(pVnode->pTq, pMsg);
L
Liu Jicong 已提交
559
    case TDMT_VND_STREAM_RECOVER_NONBLOCKING_STAGE:
L
Liu Jicong 已提交
560 561 562 563 564
      return tqProcessTaskRecover1Req(pVnode->pTq, pMsg);
    case TDMT_STREAM_RECOVER_FINISH:
      return tqProcessTaskRecoverFinishReq(pVnode->pTq, pMsg);
    case TDMT_STREAM_RECOVER_FINISH_RSP:
      return tqProcessTaskRecoverFinishRsp(pVnode->pTq, pMsg);
H
Hongze Cheng 已提交
565 566
    default:
      vError("unknown msg type:%d in fetch queue", pMsg->msgType);
567
      return TSDB_CODE_APP_ERROR;
H
Hongze Cheng 已提交
568 569 570 571 572 573
  }
}

// TODO: remove the function
void smaHandleRes(void *pVnode, int64_t smaId, const SArray *data) {
  // TODO
C
Cary Xu 已提交
574
  // blockDebugShowDataBlocks(data, __func__);
575
  tdProcessTSmaInsert(((SVnode *)pVnode)->pSma, smaId, (const char *)data);
H
Hongze Cheng 已提交
576 577
}

D
dapan1121 已提交
578
void vnodeUpdateMetaRsp(SVnode *pVnode, STableMetaRsp *pMetaRsp) {
579 580 581
  if (NULL == pMetaRsp) {
    return;
  }
L
Liu Jicong 已提交
582

D
dapan1121 已提交
583 584 585 586 587 588
  strcpy(pMetaRsp->dbFName, pVnode->config.dbname);
  pMetaRsp->dbId = pVnode->config.dbId;
  pMetaRsp->vgId = TD_VID(pVnode);
  pMetaRsp->precision = pVnode->config.tsdbCfg.precision;
}

H
Hongze Cheng 已提交
589
extern int32_t vnodeAsyncRentention(SVnode *pVnode, int64_t now);
S
Shengliang Guan 已提交
590
static int32_t vnodeProcessTrimReq(SVnode *pVnode, int64_t version, void *pReq, int32_t len, SRpcMsg *pRsp) {
H
Hongze Cheng 已提交
591
  int32_t     code = 0;
S
Shengliang Guan 已提交
592 593
  SVTrimDbReq trimReq = {0};

H
Hongze Cheng 已提交
594 595 596 597
  // decode
  if (tDeserializeSVTrimDbReq(pReq, len, &trimReq) != 0) {
    code = TSDB_CODE_INVALID_MSG;
    goto _exit;
S
Shengliang Guan 已提交
598 599
  }

C
Cary Xu 已提交
600 601
  vInfo("vgId:%d, trim vnode request will be processed, time:%d", pVnode->config.vgId, trimReq.timestamp);

H
Hongze Cheng 已提交
602 603
// process
#if 0
H
Hongze Cheng 已提交
604 605 606
  code = tsdbDoRetention(pVnode->pTsdb, trimReq.timestamp);
  if (code) goto _exit;

C
Cary Xu 已提交
607 608
  code = smaDoRetention(pVnode->pSma, trimReq.timestamp);
  if (code) goto _exit;
H
Hongze Cheng 已提交
609 610 611
#else
  vnodeAsyncRentention(pVnode, trimReq.timestamp);
#endif
C
Cary Xu 已提交
612

H
Hongze Cheng 已提交
613 614
_exit:
  return code;
S
Shengliang Guan 已提交
615 616
}

L
Liu Jicong 已提交
617
static int32_t vnodeProcessDropTtlTbReq(SVnode *pVnode, int64_t version, void *pReq, int32_t len, SRpcMsg *pRsp) {
618 619 620
  SArray *tbUids = taosArrayInit(8, sizeof(int64_t));
  if (tbUids == NULL) return TSDB_CODE_OUT_OF_MEMORY;

S
Shengliang Guan 已提交
621 622 623 624 625 626
  SVDropTtlTableReq ttlReq = {0};
  if (tDeserializeSVDropTtlTableReq(pReq, len, &ttlReq) != 0) {
    terrno = TSDB_CODE_INVALID_MSG;
    goto end;
  }

627
  vDebug("vgId:%d, drop ttl table req will be processed, time:%d", pVnode->config.vgId, ttlReq.timestamp);
S
Shengliang Guan 已提交
628
  int32_t ret = metaTtlDropTable(pVnode->pMeta, ttlReq.timestamp, tbUids);
L
Liu Jicong 已提交
629
  if (ret != 0) {
630 631
    goto end;
  }
L
Liu Jicong 已提交
632
  if (taosArrayGetSize(tbUids) > 0) {
wmmhello's avatar
wmmhello 已提交
633 634
    tqUpdateTbUidList(pVnode->pTq, tbUids, false);
  }
635

H
Hongze Cheng 已提交
636
#if 0
H
Hongze Cheng 已提交
637 638 639 640 641 642
  // process
  ret = tsdbDoRetention(pVnode->pTsdb, ttlReq.timestamp);
  if (ret) goto end;

  ret = smaDoRetention(pVnode->pSma, ttlReq.timestamp);
  if (ret) goto end;
H
Hongze Cheng 已提交
643 644
#else
  vnodeAsyncRentention(pVnode, ttlReq.timestamp);
H
Hongze Cheng 已提交
645
#endif
H
Hongze Cheng 已提交
646

647 648 649 650 651
end:
  taosArrayDestroy(tbUids);
  return ret;
}

652
static int32_t vnodeProcessCreateStbReq(SVnode *pVnode, int64_t version, void *pReq, int32_t len, SRpcMsg *pRsp) {
H
Hongze Cheng 已提交
653
  SVCreateStbReq req = {0};
H
Hongze Cheng 已提交
654
  SDecoder       coder;
H
Hongze Cheng 已提交
655

H
Hongze Cheng 已提交
656 657 658 659 660 661
  pRsp->msgType = TDMT_VND_CREATE_STB_RSP;
  pRsp->code = TSDB_CODE_SUCCESS;
  pRsp->pCont = NULL;
  pRsp->contLen = 0;

  // decode and process req
H
Hongze Cheng 已提交
662
  tDecoderInit(&coder, pReq, len);
H
Hongze Cheng 已提交
663 664

  if (tDecodeSVCreateStbReq(&coder, &req) < 0) {
H
Hongze Cheng 已提交
665 666
    pRsp->code = terrno;
    goto _err;
H
Hongze Cheng 已提交
667 668
  }

H
Hongze Cheng 已提交
669
  if (metaCreateSTable(pVnode->pMeta, version, &req) < 0) {
H
Hongze Cheng 已提交
670 671
    pRsp->code = terrno;
    goto _err;
H
Hongze Cheng 已提交
672 673
  }

674
  if (tdProcessRSmaCreate(pVnode->pSma, &req) < 0) {
675 676 677
    pRsp->code = terrno;
    goto _err;
  }
C
Cary Xu 已提交
678

H
Hongze Cheng 已提交
679
  tDecoderClear(&coder);
H
Hongze Cheng 已提交
680
  return 0;
H
Hongze Cheng 已提交
681 682

_err:
H
Hongze Cheng 已提交
683
  tDecoderClear(&coder);
H
Hongze Cheng 已提交
684
  return -1;
H
Hongze Cheng 已提交
685 686
}

687
static int32_t vnodeProcessCreateTbReq(SVnode *pVnode, int64_t version, void *pReq, int32_t len, SRpcMsg *pRsp) {
H
Hongze Cheng 已提交
688
  SDecoder           decoder = {0};
689
  SEncoder           encoder = {0};
690
  int32_t            rcode = 0;
H
Hongze Cheng 已提交
691 692
  SVCreateTbBatchReq req = {0};
  SVCreateTbReq     *pCreateReq;
H
Hongze Cheng 已提交
693 694 695
  SVCreateTbBatchRsp rsp = {0};
  SVCreateTbRsp      cRsp = {0};
  char               tbName[TSDB_TABLE_FNAME_LEN];
C
Cary Xu 已提交
696
  STbUidStore       *pStore = NULL;
697
  SArray            *tbUids = NULL;
H
Hongze Cheng 已提交
698 699

  pRsp->msgType = TDMT_VND_CREATE_TABLE_RSP;
H
Hongze Cheng 已提交
700 701 702
  pRsp->code = TSDB_CODE_SUCCESS;
  pRsp->pCont = NULL;
  pRsp->contLen = 0;
H
Hongze Cheng 已提交
703

H
Hongze Cheng 已提交
704
  // decode
H
Hongze Cheng 已提交
705 706
  tDecoderInit(&decoder, pReq, len);
  if (tDecodeSVCreateTbBatchReq(&decoder, &req) < 0) {
H
Hongze Cheng 已提交
707 708 709 710
    rcode = -1;
    terrno = TSDB_CODE_INVALID_MSG;
    goto _exit;
  }
H
Hongze Cheng 已提交
711

H
Hongze Cheng 已提交
712
  rsp.pArray = taosArrayInit(req.nReqs, sizeof(cRsp));
713 714
  tbUids = taosArrayInit(req.nReqs, sizeof(int64_t));
  if (rsp.pArray == NULL || tbUids == NULL) {
H
Hongze Cheng 已提交
715 716 717 718 719
    rcode = -1;
    terrno = TSDB_CODE_OUT_OF_MEMORY;
    goto _exit;
  }

H
Hongze Cheng 已提交
720
  // loop to create table
721
  for (int32_t iReq = 0; iReq < req.nReqs; iReq++) {
H
Hongze Cheng 已提交
722
    pCreateReq = req.pReqs + iReq;
D
dapan1121 已提交
723
    memset(&cRsp, 0, sizeof(cRsp));
H
Hongze Cheng 已提交
724

C
Cary Xu 已提交
725 726 727 728
    if ((terrno = grantCheck(TSDB_GRANT_TIMESERIES)) < 0) {
      rcode = -1;
      goto _exit;
    }
L
Liu Jicong 已提交
729

wafwerar's avatar
wafwerar 已提交
730 731 732 733 734
    if ((terrno = grantCheck(TSDB_GRANT_TABLE)) < 0) {
      rcode = -1;
      goto _exit;
    }

H
Hongze Cheng 已提交
735 736 737 738 739 740 741 742 743
    // validate hash
    sprintf(tbName, "%s.%s", pVnode->config.dbname, pCreateReq->name);
    if (vnodeValidateTableHash(pVnode, tbName) < 0) {
      cRsp.code = TSDB_CODE_VND_HASH_MISMATCH;
      taosArrayPush(rsp.pArray, &cRsp);
      continue;
    }

    // do create table
744
    if (metaCreateTable(pVnode->pMeta, version, pCreateReq, &cRsp.pMeta) < 0) {
H
Hongze Cheng 已提交
745 746 747 748 749
      if (pCreateReq->flags & TD_CREATE_IF_NOT_EXISTS && terrno == TSDB_CODE_TDB_TABLE_ALREADY_EXIST) {
        cRsp.code = TSDB_CODE_SUCCESS;
      } else {
        cRsp.code = terrno;
      }
H
Hongze Cheng 已提交
750
    } else {
H
Hongze Cheng 已提交
751
      cRsp.code = TSDB_CODE_SUCCESS;
752
      tdFetchTbUidList(pVnode->pSma, &pStore, pCreateReq->ctb.suid, pCreateReq->uid);
753
      taosArrayPush(tbUids, &pCreateReq->uid);
L
Liu Jicong 已提交
754
      vnodeUpdateMetaRsp(pVnode, cRsp.pMeta);
H
Hongze Cheng 已提交
755
    }
H
Hongze Cheng 已提交
756 757

    taosArrayPush(rsp.pArray, &cRsp);
H
Hongze Cheng 已提交
758 759
  }

H
Haojun Liao 已提交
760
  vDebug("vgId:%d, add %d new created tables into query table list", TD_VID(pVnode), (int32_t)taosArrayGetSize(tbUids));
761
  tqUpdateTbUidList(pVnode->pTq, tbUids, true);
762
  if (tdUpdateTbUidList(pVnode->pSma, pStore, true) < 0) {
C
Cary Xu 已提交
763 764
    goto _exit;
  }
765
  tdUidStoreFree(pStore);
C
Cary Xu 已提交
766

H
Hongze Cheng 已提交
767
  // prepare rsp
768
  int32_t ret = 0;
wafwerar's avatar
wafwerar 已提交
769
  tEncodeSize(tEncodeSVCreateTbBatchRsp, &rsp, pRsp->contLen, ret);
H
Hongze Cheng 已提交
770 771 772 773 774 775
  pRsp->pCont = rpcMallocCont(pRsp->contLen);
  if (pRsp->pCont == NULL) {
    terrno = TSDB_CODE_OUT_OF_MEMORY;
    rcode = -1;
    goto _exit;
  }
H
Hongze Cheng 已提交
776 777
  tEncoderInit(&encoder, pRsp->pCont, pRsp->contLen);
  tEncodeSVCreateTbBatchRsp(&encoder, &rsp);
H
Hongze Cheng 已提交
778

H
Hongze Cheng 已提交
779
_exit:
wmmhello's avatar
wmmhello 已提交
780 781
  for (int32_t iReq = 0; iReq < req.nReqs; iReq++) {
    pCreateReq = req.pReqs + iReq;
782
    taosMemoryFree(pCreateReq->comment);
wmmhello's avatar
wmmhello 已提交
783 784
    taosArrayDestroy(pCreateReq->ctb.tagName);
  }
785
  taosArrayDestroyEx(rsp.pArray, tFreeSVCreateTbRsp);
786
  taosArrayDestroy(tbUids);
H
Hongze Cheng 已提交
787 788
  tDecoderClear(&decoder);
  tEncoderClear(&encoder);
H
Hongze Cheng 已提交
789
  return rcode;
H
Hongze Cheng 已提交
790 791
}

792
static int32_t vnodeProcessAlterStbReq(SVnode *pVnode, int64_t version, void *pReq, int32_t len, SRpcMsg *pRsp) {
H
Hongze Cheng 已提交
793 794 795 796 797 798 799 800 801 802 803 804 805 806 807
  SVCreateStbReq req = {0};
  SDecoder       dc = {0};

  pRsp->msgType = TDMT_VND_ALTER_STB_RSP;
  pRsp->code = TSDB_CODE_SUCCESS;
  pRsp->pCont = NULL;
  pRsp->contLen = 0;

  tDecoderInit(&dc, pReq, len);

  // decode req
  if (tDecodeSVCreateStbReq(&dc, &req) < 0) {
    terrno = TSDB_CODE_INVALID_MSG;
    tDecoderClear(&dc);
    return -1;
H
Hongze Cheng 已提交
808
  }
H
Hongze Cheng 已提交
809 810 811 812 813 814 815 816 817

  if (metaAlterSTable(pVnode->pMeta, version, &req) < 0) {
    pRsp->code = terrno;
    tDecoderClear(&dc);
    return -1;
  }

  tDecoderClear(&dc);

H
Hongze Cheng 已提交
818 819 820
  return 0;
}

821
static int32_t vnodeProcessDropStbReq(SVnode *pVnode, int64_t version, void *pReq, int32_t len, SRpcMsg *pRsp) {
H
Hongze Cheng 已提交
822
  SVDropStbReq req = {0};
823
  int32_t      rcode = TSDB_CODE_SUCCESS;
H
Hongze Cheng 已提交
824
  SDecoder     decoder = {0};
825
  SArray      *tbUidList = NULL;
H
Hongze Cheng 已提交
826 827 828 829 830 831

  pRsp->msgType = TDMT_VND_CREATE_STB_RSP;
  pRsp->pCont = NULL;
  pRsp->contLen = 0;

  // decode request
H
Hongze Cheng 已提交
832 833
  tDecoderInit(&decoder, pReq, len);
  if (tDecodeSVDropStbReq(&decoder, &req) < 0) {
H
Hongze Cheng 已提交
834 835 836 837 838
    rcode = TSDB_CODE_INVALID_MSG;
    goto _exit;
  }

  // process request
839 840 841 842 843 844 845 846
  tbUidList = taosArrayInit(8, sizeof(int64_t));
  if (tbUidList == NULL) goto _exit;
  if (metaDropSTable(pVnode->pMeta, version, &req, tbUidList) < 0) {
    rcode = terrno;
    goto _exit;
  }

  if (tqUpdateTbUidList(pVnode->pTq, tbUidList, false) < 0) {
847 848 849
    rcode = terrno;
    goto _exit;
  }
H
Hongze Cheng 已提交
850

851 852 853 854 855
  if (tdProcessRSmaDrop(pVnode->pSma, &req) < 0) {
    rcode = terrno;
    goto _exit;
  }

H
Hongze Cheng 已提交
856 857
  // return rsp
_exit:
858
  if (tbUidList) taosArrayDestroy(tbUidList);
H
Hongze Cheng 已提交
859
  pRsp->code = rcode;
H
Hongze Cheng 已提交
860
  tDecoderClear(&decoder);
H
Hongze Cheng 已提交
861 862 863
  return 0;
}

864
static int32_t vnodeProcessAlterTbReq(SVnode *pVnode, int64_t version, void *pReq, int32_t len, SRpcMsg *pRsp) {
D
dapan1121 已提交
865 866 867
  SVAlterTbReq  vAlterTbReq = {0};
  SVAlterTbRsp  vAlterTbRsp = {0};
  SDecoder      dc = {0};
868 869
  int32_t       rcode = 0;
  int32_t       ret;
D
dapan1121 已提交
870 871
  SEncoder      ec = {0};
  STableMetaRsp vMetaRsp = {0};
H
Hongze Cheng 已提交
872 873 874 875 876 877 878 879 880 881

  pRsp->msgType = TDMT_VND_ALTER_TABLE_RSP;
  pRsp->pCont = NULL;
  pRsp->contLen = 0;
  pRsp->code = TSDB_CODE_SUCCESS;

  tDecoderInit(&dc, pReq, len);

  // decode
  if (tDecodeSVAlterTbReq(&dc, &vAlterTbReq) < 0) {
H
Hongze Cheng 已提交
882
    vAlterTbRsp.code = TSDB_CODE_INVALID_MSG;
H
Hongze Cheng 已提交
883
    tDecoderClear(&dc);
H
Hongze Cheng 已提交
884 885
    rcode = -1;
    goto _exit;
H
Hongze Cheng 已提交
886 887 888
  }

  // process
D
dapan1121 已提交
889
  if (metaAlterTable(pVnode->pMeta, version, &vAlterTbReq, &vMetaRsp) < 0) {
wmmhello's avatar
wmmhello 已提交
890
    vAlterTbRsp.code = terrno;
H
Hongze Cheng 已提交
891
    tDecoderClear(&dc);
H
Hongze Cheng 已提交
892 893
    rcode = -1;
    goto _exit;
H
Hongze Cheng 已提交
894 895
  }
  tDecoderClear(&dc);
H
Hongze Cheng 已提交
896

D
dapan1121 已提交
897 898 899 900 901
  if (NULL != vMetaRsp.pSchemas) {
    vnodeUpdateMetaRsp(pVnode, &vMetaRsp);
    vAlterTbRsp.pMeta = &vMetaRsp;
  }

H
Hongze Cheng 已提交
902 903 904 905 906 907
_exit:
  tEncodeSize(tEncodeSVAlterTbRsp, &vAlterTbRsp, pRsp->contLen, ret);
  pRsp->pCont = rpcMallocCont(pRsp->contLen);
  tEncoderInit(&ec, pRsp->pCont, pRsp->contLen);
  tEncodeSVAlterTbRsp(&ec, &vAlterTbRsp);
  tEncoderClear(&ec);
M
Minglei Jin 已提交
908 909 910
  if (vMetaRsp.pSchemas) {
    taosMemoryFree(vMetaRsp.pSchemas);
  }
H
Hongze Cheng 已提交
911 912 913
  return 0;
}

914
static int32_t vnodeProcessDropTbReq(SVnode *pVnode, int64_t version, void *pReq, int32_t len, SRpcMsg *pRsp) {
H
Hongze Cheng 已提交
915 916
  SVDropTbBatchReq req = {0};
  SVDropTbBatchRsp rsp = {0};
H
Hongze Cheng 已提交
917
  SDecoder         decoder = {0};
H
Hongze Cheng 已提交
918
  SEncoder         encoder = {0};
919
  int32_t          ret;
920
  SArray          *tbUids = NULL;
921
  STbUidStore     *pStore = NULL;
H
Hongze Cheng 已提交
922

H
Hongze Cheng 已提交
923
  pRsp->msgType = TDMT_VND_DROP_TABLE_RSP;
H
Hongze Cheng 已提交
924 925 926
  pRsp->pCont = NULL;
  pRsp->contLen = 0;
  pRsp->code = TSDB_CODE_SUCCESS;
H
Hongze Cheng 已提交
927 928

  // decode req
H
Hongze Cheng 已提交
929 930
  tDecoderInit(&decoder, pReq, len);
  ret = tDecodeSVDropTbBatchReq(&decoder, &req);
H
Hongze Cheng 已提交
931 932 933 934 935
  if (ret < 0) {
    terrno = TSDB_CODE_INVALID_MSG;
    pRsp->code = terrno;
    goto _exit;
  }
H
Hongze Cheng 已提交
936 937

  // process req
938
  tbUids = taosArrayInit(req.nReqs, sizeof(int64_t));
H
Hongze Cheng 已提交
939
  rsp.pArray = taosArrayInit(req.nReqs, sizeof(SVDropTbRsp));
940 941
  if (tbUids == NULL || rsp.pArray == NULL) goto _exit;

942
  for (int32_t iReq = 0; iReq < req.nReqs; iReq++) {
H
Hongze Cheng 已提交
943 944
    SVDropTbReq *pDropTbReq = req.pReqs + iReq;
    SVDropTbRsp  dropTbRsp = {0};
945
    tb_uid_t     tbUid = 0;
H
Hongze Cheng 已提交
946

H
Hongze Cheng 已提交
947
    /* code */
948
    ret = metaDropTable(pVnode->pMeta, version, pDropTbReq, tbUids, &tbUid);
H
Hongze Cheng 已提交
949
    if (ret < 0) {
950
      if (pDropTbReq->igNotExists && terrno == TSDB_CODE_TDB_TABLE_NOT_EXIST) {
H
Hongze Cheng 已提交
951 952 953 954
        dropTbRsp.code = TSDB_CODE_SUCCESS;
      } else {
        dropTbRsp.code = terrno;
      }
H
Hongze Cheng 已提交
955
    } else {
H
Hongze Cheng 已提交
956
      dropTbRsp.code = TSDB_CODE_SUCCESS;
957
      if (tbUid > 0) tdFetchTbUidList(pVnode->pSma, &pStore, pDropTbReq->suid, tbUid);
H
Hongze Cheng 已提交
958 959 960 961 962
    }

    taosArrayPush(rsp.pArray, &dropTbRsp);
  }

963
  tqUpdateTbUidList(pVnode->pTq, tbUids, false);
964
  tdUpdateTbUidList(pVnode->pSma, pStore, false);
965

H
Hongze Cheng 已提交
966
_exit:
967
  taosArrayDestroy(tbUids);
968
  tdUidStoreFree(pStore);
H
Hongze Cheng 已提交
969
  tDecoderClear(&decoder);
H
Hongze Cheng 已提交
970 971 972 973 974
  tEncodeSize(tEncodeSVDropTbBatchRsp, &rsp, pRsp->contLen, ret);
  pRsp->pCont = rpcMallocCont(pRsp->contLen);
  tEncoderInit(&encoder, pRsp->pCont, pRsp->contLen);
  tEncodeSVDropTbBatchRsp(&encoder, &rsp);
  tEncoderClear(&encoder);
975
  taosArrayDestroy(rsp.pArray);
H
Hongze Cheng 已提交
976 977 978
  return 0;
}

979 980
static int32_t vnodeDebugPrintSingleSubmitMsg(SMeta *pMeta, SSubmitBlk *pBlock, SSubmitMsgIter *msgIter,
                                              const char *tags) {
D
dapan 已提交
981 982 983 984
  SSubmitBlkIter blkIter = {0};
  STSchema      *pSchema = NULL;
  tb_uid_t       suid = 0;
  STSRow        *row = NULL;
C
Cary Xu 已提交
985
  int32_t        rv = -1;
D
dapan 已提交
986 987 988

  tInitSubmitBlkIter(msgIter, pBlock, &blkIter);
  if (blkIter.row == NULL) return 0;
989 990 991 992 993

  pSchema = metaGetTbTSchema(pMeta, msgIter->suid, TD_ROW_SVER(blkIter.row), 1);  // TODO: use the real schema
  if (pSchema) {
    suid = msgIter->suid;
    rv = TD_ROW_SVER(blkIter.row);
D
dapan 已提交
994 995 996 997 998 999 1000 1001 1002 1003 1004 1005 1006 1007 1008 1009
  }
  if (!pSchema) {
    printf("%s:%d no valid schema\n", tags, __LINE__);
    return -1;
  }
  char __tags[128] = {0};
  snprintf(__tags, 128, "%s: uid %" PRIi64 " ", tags, msgIter->uid);
  while ((row = tGetSubmitBlkNext(&blkIter))) {
    tdSRowPrint(row, pSchema, __tags);
  }

  taosMemoryFreeClear(pSchema);

  return TSDB_CODE_SUCCESS;
}

1010 1011 1012 1013 1014 1015 1016 1017 1018 1019 1020 1021 1022
typedef struct SSubmitReqConvertCxt {
  SSubmitMsgIter msgIter;
  SSubmitBlk    *pBlock;
  SSubmitBlkIter blkIter;
  STSRow        *pRow;
  STSRowIter     rowIter;
  SSubmitTbData *pTbData;
  STSchema      *pTbSchema;
  SArray        *pColValues;
} SSubmitReqConvertCxt;

static int32_t vnodeResetTableCxt(SMeta *pMeta, SSubmitReqConvertCxt *pCxt) {
  taosMemoryFreeClear(pCxt->pTbSchema);
1023 1024
  pCxt->pTbSchema = metaGetTbTSchema(pMeta, pCxt->msgIter.suid > 0 ? pCxt->msgIter.suid : pCxt->msgIter.uid,
                                     pCxt->msgIter.sversion, 1);
1025 1026 1027 1028 1029 1030 1031 1032 1033 1034 1035 1036 1037 1038 1039 1040 1041 1042 1043 1044 1045 1046 1047 1048 1049 1050 1051 1052 1053 1054 1055 1056 1057 1058 1059 1060 1061 1062 1063 1064 1065 1066 1067 1068 1069 1070 1071 1072 1073 1074 1075 1076 1077 1078 1079 1080 1081 1082 1083 1084 1085 1086 1087 1088 1089 1090 1091 1092 1093 1094 1095 1096 1097 1098 1099 1100 1101 1102 1103 1104 1105 1106 1107 1108 1109 1110 1111 1112 1113 1114 1115 1116 1117 1118 1119 1120 1121 1122 1123 1124 1125 1126 1127 1128 1129 1130 1131 1132 1133 1134 1135 1136 1137 1138 1139 1140 1141 1142 1143 1144 1145 1146 1147 1148 1149 1150 1151 1152 1153 1154 1155 1156 1157 1158 1159 1160 1161 1162 1163 1164 1165 1166 1167 1168 1169 1170 1171 1172 1173 1174 1175 1176 1177 1178 1179 1180 1181 1182 1183 1184 1185 1186 1187 1188 1189 1190
  if (NULL == pCxt->pTbSchema) {
    return TSDB_CODE_INVALID_MSG;
  }
  tdSTSRowIterInit(&pCxt->rowIter, pCxt->pTbSchema);

  tDestroySSubmitTbData(pCxt->pTbData, TSDB_MSG_FLG_ENCODE);
  if (NULL == pCxt->pTbData) {
    pCxt->pTbData = taosMemoryCalloc(1, sizeof(SSubmitTbData));
    if (NULL == pCxt->pTbData) {
      return TSDB_CODE_OUT_OF_MEMORY;
    }
  }
  pCxt->pTbData->flags = 0;
  pCxt->pTbData->suid = pCxt->msgIter.suid;
  pCxt->pTbData->uid = pCxt->msgIter.uid;
  pCxt->pTbData->sver = pCxt->msgIter.sversion;
  pCxt->pTbData->pCreateTbReq = NULL;
  pCxt->pTbData->aRowP = taosArrayInit(128, POINTER_BYTES);
  if (NULL == pCxt->pTbData->aRowP) {
    return TSDB_CODE_OUT_OF_MEMORY;
  }

  taosArrayDestroy(pCxt->pColValues);
  pCxt->pColValues = taosArrayInit(pCxt->pTbSchema->numOfCols, sizeof(SColVal));
  if (NULL == pCxt->pColValues) {
    return TSDB_CODE_OUT_OF_MEMORY;
  }
  for (int32_t i = 0; i < pCxt->pTbSchema->numOfCols; ++i) {
    SColVal val = COL_VAL_NONE(pCxt->pTbSchema->columns[i].colId, pCxt->pTbSchema->columns[i].type);
    taosArrayPush(pCxt->pColValues, &val);
  }

  return TSDB_CODE_SUCCESS;
}

static void vnodeDestroySubmitReqConvertCxt(SSubmitReqConvertCxt *pCxt) {
  taosMemoryFreeClear(pCxt->pTbSchema);
  tDestroySSubmitTbData(pCxt->pTbData, TSDB_MSG_FLG_ENCODE);
  taosMemoryFreeClear(pCxt->pTbData);
  taosArrayDestroy(pCxt->pColValues);
}

static int32_t vnodeCellValConvertToColVal(STColumn *pCol, SCellVal *pCellVal, SColVal *pColVal) {
  if (tdValTypeIsNone(pCellVal->valType)) {
    pColVal->flag = CV_FLAG_NONE;
    return TSDB_CODE_SUCCESS;
  }

  if (tdValTypeIsNull(pCellVal->valType)) {
    pColVal->flag = CV_FLAG_NULL;
    return TSDB_CODE_SUCCESS;
  }

  if (IS_VAR_DATA_TYPE(pCol->type)) {
    pColVal->value.nData = varDataLen(pCellVal->val);
    pColVal->value.pData = varDataVal(pCellVal->val);
  } else if (TSDB_DATA_TYPE_FLOAT == pCol->type) {
    float f = GET_FLOAT_VAL(pCellVal->val);
    memcpy(&pColVal->value.val, &f, sizeof(f));
  } else if (TSDB_DATA_TYPE_DOUBLE == pCol->type) {
    pColVal->value.val = *(int64_t *)pCellVal->val;
  } else {
    GET_TYPED_DATA(pColVal->value.val, int64_t, pCol->type, pCellVal->val);
  }

  pColVal->flag = CV_FLAG_VALUE;
  return TSDB_CODE_SUCCESS;
}

static int32_t vnodeTSRowConvertToColValArray(SSubmitReqConvertCxt *pCxt) {
  int32_t code = TSDB_CODE_SUCCESS;
  tdSTSRowIterReset(&pCxt->rowIter, pCxt->pRow);
  for (int32_t i = 0; TSDB_CODE_SUCCESS == code && i < pCxt->pTbSchema->numOfCols; ++i) {
    STColumn *pCol = pCxt->pTbSchema->columns + i;
    SCellVal  cellVal = {0};
    if (!tdSTSRowIterFetch(&pCxt->rowIter, pCol->colId, pCol->type, &cellVal)) {
      break;
    }
    code = vnodeCellValConvertToColVal(pCol, &cellVal, (SColVal *)taosArrayGet(pCxt->pColValues, i));
  }
  return code;
}

static int32_t vnodeDecodeCreateTbReq(SSubmitReqConvertCxt *pCxt) {
  if (pCxt->msgIter.schemaLen <= 0) {
    return TSDB_CODE_SUCCESS;
  }

  pCxt->pTbData->pCreateTbReq = taosMemoryCalloc(1, sizeof(SVCreateTbReq));
  if (NULL == pCxt->pTbData->pCreateTbReq) {
    return TSDB_CODE_OUT_OF_MEMORY;
  }

  SDecoder decoder = {0};
  tDecoderInit(&decoder, pCxt->pBlock->data, pCxt->msgIter.schemaLen);
  int32_t code = tDecodeSVCreateTbReq(&decoder, pCxt->pTbData->pCreateTbReq);
  tDecoderClear(&decoder);

  return code;
}

static int32_t vnodeSubmitReqConvertToSubmitReq2(SVnode *pVnode, SSubmitReq *pReq, SSubmitReq2 *pReq2) {
  pReq2->aSubmitTbData = taosArrayInit(128, sizeof(SSubmitTbData));
  if (NULL == pReq2->aSubmitTbData) {
    return TSDB_CODE_OUT_OF_MEMORY;
  }

  SSubmitReqConvertCxt cxt = {0};

  int32_t code = tInitSubmitMsgIter(pReq, &cxt.msgIter);
  while (TSDB_CODE_SUCCESS == code) {
    code = tGetSubmitMsgNext(&cxt.msgIter, &cxt.pBlock);
    if (TSDB_CODE_SUCCESS == code) {
      if (NULL == cxt.pBlock) {
        break;
      }
      code = vnodeResetTableCxt(pVnode->pMeta, &cxt);
    }
    if (TSDB_CODE_SUCCESS == code) {
      code = tInitSubmitBlkIter(&cxt.msgIter, cxt.pBlock, &cxt.blkIter);
    }
    if (TSDB_CODE_SUCCESS == code) {
      code = vnodeDecodeCreateTbReq(&cxt);
    }
    while (TSDB_CODE_SUCCESS == code && (cxt.pRow = tGetSubmitBlkNext(&cxt.blkIter)) != NULL) {
      code = vnodeTSRowConvertToColValArray(&cxt);
      if (TSDB_CODE_SUCCESS == code) {
        SRow **pNewRow = taosArrayReserve(cxt.pTbData->aRowP, 1);
        code = tRowBuild(cxt.pColValues, cxt.pTbSchema, pNewRow);
      }
    }
    if (TSDB_CODE_SUCCESS == code) {
      code = (NULL == taosArrayPush(pReq2->aSubmitTbData, cxt.pTbData) ? TSDB_CODE_OUT_OF_MEMORY : TSDB_CODE_SUCCESS);
    }
    if (TSDB_CODE_SUCCESS == code) {
      taosMemoryFreeClear(cxt.pTbData);
    }
  }

  vnodeDestroySubmitReqConvertCxt(&cxt);
  return code;
}

static int32_t vnodeRebuildSubmitReqMsg(SSubmitReq2 *pSubmitReq, void **ppMsg) {
  int32_t  code = TSDB_CODE_SUCCESS;
  char    *pMsg = NULL;
  uint32_t msglen = 0;
  tEncodeSize(tEncodeSSubmitReq2, pSubmitReq, msglen, code);
  if (TSDB_CODE_SUCCESS == code) {
    pMsg = taosMemoryMalloc(msglen);
    if (NULL == pMsg) {
      code = TSDB_CODE_OUT_OF_MEMORY;
    }
  }
  if (TSDB_CODE_SUCCESS == code) {
    SEncoder encoder;
    tEncoderInit(&encoder, pMsg, msglen);
    code = tEncodeSSubmitReq2(&encoder, pSubmitReq);
    tEncoderClear(&encoder);
  }
  if (TSDB_CODE_SUCCESS == code) {
    *ppMsg = pMsg;
  }
  return code;
}

1191
static int32_t vnodeProcessSubmitReq(SVnode *pVnode, int64_t version, void *pReq, int32_t len, SRpcMsg *pRsp) {
H
Hongze Cheng 已提交
1192
  int32_t code = 0;
D
dapan1121 已提交
1193
  terrno = 0;
C
Cary Xu 已提交
1194

H
Hongze Cheng 已提交
1195
  SSubmitReq2 *pSubmitReq = &(SSubmitReq2){0};
H
Hongze Cheng 已提交
1196
  SSubmitRsp2 *pSubmitRsp = &(SSubmitRsp2){0};
H
Hongze Cheng 已提交
1197
  SArray      *newTbUids = NULL;
H
Hongze Cheng 已提交
1198 1199 1200 1201
  int32_t      ret;
  SEncoder     ec = {0};

  pRsp->code = TSDB_CODE_SUCCESS;
H
Hongze Cheng 已提交
1202

1203
  void           *pAllocMsg = NULL;
1204 1205 1206 1207 1208 1209
  SSubmitReq2Msg *pMsg = (SSubmitReq2Msg *)pReq;
  if (0 == pMsg->version) {
    code = vnodeSubmitReqConvertToSubmitReq2(pVnode, (SSubmitReq *)pMsg, pSubmitReq);
    if (TSDB_CODE_SUCCESS == code) {
      code = vnodeRebuildSubmitReqMsg(pSubmitReq, &pReq);
    }
1210 1211 1212
    if (TSDB_CODE_SUCCESS == code) {
      pAllocMsg = pReq;
    }
1213 1214 1215 1216 1217 1218 1219 1220 1221 1222 1223 1224 1225 1226
    if (TSDB_CODE_SUCCESS != code) {
      goto _exit;
    }
  } else {
    // decode
    pReq = POINTER_SHIFT(pReq, sizeof(SSubmitReq2Msg));
    len -= sizeof(SSubmitReq2Msg);
    SDecoder dc = {0};
    tDecoderInit(&dc, pReq, len);
    if (tDecodeSSubmitReq2(&dc, pSubmitReq) < 0) {
      code = TSDB_CODE_INVALID_MSG;
      goto _exit;
    }
    tDecoderClear(&dc);
D
dapan 已提交
1227
  }
C
Cary Xu 已提交
1228

H
Hongze Cheng 已提交
1229
  for (int32_t i = 0; i < TARRAY_SIZE(pSubmitReq->aSubmitTbData); ++i) {
H
Hongze Cheng 已提交
1230 1231 1232 1233 1234 1235 1236 1237 1238 1239
    SSubmitTbData *pSubmitTbData = taosArrayGet(pSubmitReq->aSubmitTbData, i);

    if (pSubmitTbData->pCreateTbReq) {
      pSubmitTbData->uid = pSubmitTbData->pCreateTbReq->uid;
    } else {
      SMetaInfo info = {0};

      code = metaGetInfo(pVnode->pMeta, pSubmitTbData->uid, &info, NULL);
      if (code) {
        code = TSDB_CODE_TDB_TABLE_NOT_EXIST;
D
dapan1121 已提交
1240
        vWarn("vgId:%d, table uid:%" PRId64 " not exists", TD_VID(pVnode), pSubmitTbData->uid);
H
Hongze Cheng 已提交
1241 1242 1243 1244 1245 1246 1247 1248 1249 1250 1251 1252 1253 1254 1255 1256 1257
        goto _exit;
      }

      if (info.suid != pSubmitTbData->suid) {
        code = TSDB_CODE_INVALID_MSG;
        goto _exit;
      }

      if (info.suid) {
        metaGetInfo(pVnode->pMeta, info.suid, &info, NULL);
      }

      if (pSubmitTbData->sver != info.skmVer) {
        code = TSDB_CODE_TDB_INVALID_TABLE_SCHEMA_VER;
        goto _exit;
      }
    }
H
Hongze Cheng 已提交
1258 1259 1260 1261 1262 1263 1264 1265 1266 1267 1268 1269 1270 1271 1272 1273 1274 1275 1276 1277 1278 1279 1280

    if (pSubmitTbData->flags & SUBMIT_REQ_COLUMN_DATA_FORMAT) {
      int32_t   nColData = TARRAY_SIZE(pSubmitTbData->aCol);
      SColData *aColData = (SColData *)TARRAY_DATA(pSubmitTbData->aCol);

      if (nColData <= 0) {
        code = TSDB_CODE_INVALID_MSG;
        goto _exit;
      }

      if (aColData[0].cid != PRIMARYKEY_TIMESTAMP_COL_ID || aColData[0].type != TSDB_DATA_TYPE_TIMESTAMP ||
          aColData[0].nVal <= 0) {
        code = TSDB_CODE_INVALID_MSG;
        goto _exit;
      }

      for (int32_t i = 1; i < nColData; i++) {
        if (aColData[i].nVal != aColData[0].nVal) {
          code = TSDB_CODE_INVALID_MSG;
          goto _exit;
        }
      }
    }
H
Hongze Cheng 已提交
1281 1282
  }

L
Liu Jicong 已提交
1283 1284
  vDebug("vgId:%d, submit block size %d", TD_VID(pVnode), (int32_t)taosArrayGetSize(pSubmitReq->aSubmitTbData));

H
Hongze Cheng 已提交
1285
  // loop to handle
H
Hongze Cheng 已提交
1286
  for (int32_t i = 0; i < TARRAY_SIZE(pSubmitReq->aSubmitTbData); ++i) {
H
Hongze Cheng 已提交
1287 1288 1289 1290
    SSubmitTbData *pSubmitTbData = taosArrayGet(pSubmitReq->aSubmitTbData, i);

    // create table
    if (pSubmitTbData->pCreateTbReq) {
H
Hongze Cheng 已提交
1291 1292 1293 1294 1295 1296 1297 1298 1299 1300 1301
      // check (TODO: move check to create table)
      code = grantCheck(TSDB_GRANT_TIMESERIES);
      if (code) goto _exit;

      code = grantCheck(TSDB_GRANT_TABLE);
      if (code) goto _exit;

      // alloc if need
      if (pSubmitRsp->aCreateTbRsp == NULL &&
          (pSubmitRsp->aCreateTbRsp = taosArrayInit(TARRAY_SIZE(pSubmitReq->aSubmitTbData), sizeof(SVCreateTbRsp))) ==
              NULL) {
H
Hongze Cheng 已提交
1302
        code = TSDB_CODE_OUT_OF_MEMORY;
H
Hongze Cheng 已提交
1303 1304 1305 1306 1307 1308 1309 1310 1311 1312 1313
        goto _exit;
      }

      SVCreateTbRsp *pCreateTbRsp = taosArrayReserve(pSubmitRsp->aCreateTbRsp, 1);

      // create table
      if (metaCreateTable(pVnode->pMeta, version, pSubmitTbData->pCreateTbReq, &pCreateTbRsp->pMeta) ==
          0) {  // create table success

        if (newTbUids == NULL &&
            (newTbUids = taosArrayInit(TARRAY_SIZE(pSubmitReq->aSubmitTbData), sizeof(int64_t))) == NULL) {
H
Hongze Cheng 已提交
1314
          code = TSDB_CODE_OUT_OF_MEMORY;
H
Hongze Cheng 已提交
1315 1316 1317 1318 1319 1320
          goto _exit;
        }

        taosArrayPush(newTbUids, &pSubmitTbData->uid);

        if (pCreateTbRsp->pMeta) {
D
dapan1121 已提交
1321
          vnodeUpdateMetaRsp(pVnode, pCreateTbRsp->pMeta);
H
Hongze Cheng 已提交
1322 1323 1324 1325 1326 1327
        }
      } else {  // create table failed
        if (terrno != TSDB_CODE_TDB_TABLE_ALREADY_EXIST) {
          code = terrno;
          goto _exit;
        }
1328
        pSubmitTbData->uid = pSubmitTbData->pCreateTbReq->uid;  // update uid if table exist for using below
H
Hongze Cheng 已提交
1329
      }
H
Hongze Cheng 已提交
1330 1331 1332
    }

    // insert data
H
Hongze Cheng 已提交
1333 1334 1335 1336 1337
    int32_t affectedRows;
    code = tsdbInsertTableData(pVnode->pTsdb, version, pSubmitTbData, &affectedRows);
    if (code) goto _exit;

    pSubmitRsp->affectedRows += affectedRows;
H
Hongze Cheng 已提交
1338
  }
H
Hongze Cheng 已提交
1339

H
Hongze Cheng 已提交
1340 1341 1342 1343 1344
  // update table uid list
  if (taosArrayGetSize(newTbUids) > 0) {
    vDebug("vgId:%d, add %d table into query table list in handling submit", TD_VID(pVnode),
           (int32_t)taosArrayGetSize(newTbUids));
    tqUpdateTbUidList(pVnode->pTq, newTbUids, true);
H
Hongze Cheng 已提交
1345
  }
H
Hongze Cheng 已提交
1346 1347 1348 1349 1350 1351 1352 1353 1354 1355

_exit:
  // message
  pRsp->code = code;
  tEncodeSize(tEncodeSSubmitRsp2, pSubmitRsp, pRsp->contLen, ret);
  pRsp->pCont = rpcMallocCont(pRsp->contLen);
  tEncoderInit(&ec, pRsp->pCont, pRsp->contLen);
  tEncodeSSubmitRsp2(&ec, pSubmitRsp);
  tEncoderClear(&ec);

H
Hongze Cheng 已提交
1356 1357 1358 1359 1360 1361
  // update statistics
  atomic_add_fetch_64(&pVnode->statis.nInsert, pSubmitRsp->affectedRows);
  atomic_add_fetch_64(&pVnode->statis.nInsertSuccess, pSubmitRsp->affectedRows);
  atomic_add_fetch_64(&pVnode->statis.nBatchInsert, 1);
  if (code == 0) {
    atomic_add_fetch_64(&pVnode->statis.nBatchInsertSuccess, 1);
K
kailixu 已提交
1362
    tdProcessRSmaSubmit(pVnode->pSma, version, pSubmitReq, pReq, len, STREAM_INPUT__DATA_SUBMIT);
H
Hongze Cheng 已提交
1363 1364
  }

H
Hongze Cheng 已提交
1365 1366
  // clear
  taosArrayDestroy(newTbUids);
1367
  tDestroySSubmitReq2(pSubmitReq, 0 == pMsg->version ? TSDB_MSG_FLG_CMPT : TSDB_MSG_FLG_DECODE);
H
Hongze Cheng 已提交
1368
  tDestroySSubmitRsp2(pSubmitRsp, TSDB_MSG_FLG_ENCODE);
H
Hongze Cheng 已提交
1369

D
dapan1121 已提交
1370 1371
  if (code) terrno = code;

1372
  taosMemoryFree(pAllocMsg);
1373

H
Hongze Cheng 已提交
1374
  return code;
L
Liu Jicong 已提交
1375
}
1376

1377
static int32_t vnodeProcessCreateTSmaReq(SVnode *pVnode, int64_t version, void *pReq, int32_t len, SRpcMsg *pRsp) {
1378
  SVCreateTSmaReq req = {0};
C
Cary Xu 已提交
1379
  SDecoder        coder = {0};
1380

C
Cary Xu 已提交
1381 1382 1383 1384 1385 1386
  if (pRsp) {
    pRsp->msgType = TDMT_VND_CREATE_SMA_RSP;
    pRsp->code = TSDB_CODE_SUCCESS;
    pRsp->pCont = NULL;
    pRsp->contLen = 0;
  }
1387 1388 1389 1390 1391

  // decode and process req
  tDecoderInit(&coder, pReq, len);

  if (tDecodeSVCreateTSmaReq(&coder, &req) < 0) {
C
Cary Xu 已提交
1392 1393
    terrno = TSDB_CODE_MSG_DECODE_ERROR;
    if (pRsp) pRsp->code = terrno;
1394
    goto _err;
1395
  }
C
Cary Xu 已提交
1396

C
Cary Xu 已提交
1397
  if (tdProcessTSmaCreate(pVnode->pSma, version, (const char *)&req) < 0) {
C
Cary Xu 已提交
1398
    if (pRsp) pRsp->code = terrno;
1399
    goto _err;
1400
  }
C
Cary Xu 已提交
1401

1402
  tDecoderClear(&coder);
S
Shengliang Guan 已提交
1403
  vDebug("vgId:%d, success to create tsma %s:%" PRIi64 " version %" PRIi64 " for table %" PRIi64, TD_VID(pVnode),
C
Cary Xu 已提交
1404
         req.indexName, req.indexUid, version, req.tableUid);
H
Hongze Cheng 已提交
1405
  return 0;
1406 1407 1408

_err:
  tDecoderClear(&coder);
S
Shengliang Guan 已提交
1409
  vError("vgId:%d, failed to create tsma %s:%" PRIi64 " version %" PRIi64 "for table %" PRIi64 " since %s",
C
Cary Xu 已提交
1410
         TD_VID(pVnode), req.indexName, req.indexUid, version, req.tableUid, terrstr());
1411
  return -1;
L
Liu Jicong 已提交
1412
}
C
Cary Xu 已提交
1413 1414 1415 1416 1417 1418 1419 1420 1421 1422 1423 1424

/**
 * @brief specific for smaDstVnode
 *
 * @param pVnode
 * @param pCont
 * @param contLen
 * @return int32_t
 */
int32_t vnodeProcessCreateTSma(SVnode *pVnode, void *pCont, uint32_t contLen) {
  return vnodeProcessCreateTSmaReq(pVnode, 1, pCont, contLen, NULL);
}
1425 1426 1427 1428 1429 1430 1431 1432 1433

static int32_t vnodeProcessAlterConfirmReq(SVnode *pVnode, int64_t version, void *pReq, int32_t len, SRpcMsg *pRsp) {
  vInfo("vgId:%d, alter replica confim msg is processed", TD_VID(pVnode));
  pRsp->msgType = TDMT_VND_ALTER_CONFIRM_RSP;
  pRsp->code = TSDB_CODE_SUCCESS;
  pRsp->pCont = NULL;
  pRsp->contLen = 0;

  return 0;
1434
}
S
Shengliang Guan 已提交
1435

1436
static int32_t vnodeProcessAlterConfigReq(SVnode *pVnode, int64_t version, void *pReq, int32_t len, SRpcMsg *pRsp) {
S
Shengliang Guan 已提交
1437 1438
  bool walChanged = false;
  bool tsdbChanged = false;
1439

S
Shengliang Guan 已提交
1440 1441
  SAlterVnodeConfigReq req = {0};
  if (tDeserializeSAlterVnodeConfigReq(pReq, len, &req) != 0) {
1442 1443 1444 1445
    terrno = TSDB_CODE_INVALID_MSG;
    return TSDB_CODE_INVALID_MSG;
  }

1446
  vInfo("vgId:%d, start to alter vnode config, page:%d pageSize:%d buffer:%d szPage:%d szBuf:%" PRIu64
S
Shengliang Guan 已提交
1447
        " cacheLast:%d cacheLastSize:%d days:%d keep0:%d keep1:%d keep2:%d fsync:%d level:%d",
1448 1449
        TD_VID(pVnode), req.pages, req.pageSize, req.buffer, req.pageSize * 1024, (uint64_t)req.buffer * 1024 * 1024,
        req.cacheLast, req.cacheLastSize, req.daysPerFile, req.daysToKeep0, req.daysToKeep1, req.daysToKeep2,
S
Shengliang Guan 已提交
1450
        req.walFsyncPeriod, req.walLevel);
1451 1452 1453

  if (pVnode->config.cacheLastSize != req.cacheLastSize) {
    pVnode->config.cacheLastSize = req.cacheLastSize;
1454 1455
    tsdbCacheSetCapacity(pVnode, (size_t)pVnode->config.cacheLastSize * 1024 * 1024);
  }
1456

1457
  if (pVnode->config.szBuf != req.buffer * 1024LL * 1024LL) {
S
Shengliang Guan 已提交
1458
    vInfo("vgId:%d, vnode buffer is changed from %" PRId64 " to %" PRId64, TD_VID(pVnode), pVnode->config.szBuf,
1459
          (uint64_t)(req.buffer * 1024LL * 1024LL));
1460
    pVnode->config.szBuf = req.buffer * 1024LL * 1024LL;
H
Hongze Cheng 已提交
1461 1462
  }

1463 1464
  if (pVnode->config.szCache != req.pages) {
    if (metaAlterCache(pVnode->pMeta, req.pages) < 0) {
S
Shengliang Guan 已提交
1465
      vError("vgId:%d, failed to change vnode pages from %d to %d failed since %s", TD_VID(pVnode),
1466
             pVnode->config.szCache, req.pages, tstrerror(errno));
H
Hongze Cheng 已提交
1467 1468
      return errno;
    } else {
S
Shengliang Guan 已提交
1469
      vInfo("vgId:%d, vnode pages is changed from %d to %d", TD_VID(pVnode), pVnode->config.szCache, req.pages);
1470
      pVnode->config.szCache = req.pages;
H
Hongze Cheng 已提交
1471
    }
H
Hongze Cheng 已提交
1472 1473
  }

1474 1475
  if (pVnode->config.cacheLast != req.cacheLast) {
    pVnode->config.cacheLast = req.cacheLast;
1476 1477
  }

1478 1479
  if (pVnode->config.walCfg.fsyncPeriod != req.walFsyncPeriod) {
    pVnode->config.walCfg.fsyncPeriod = req.walFsyncPeriod;
1480 1481 1482 1483

    walChanged = true;
  }

1484 1485
  if (pVnode->config.walCfg.level != req.walLevel) {
    pVnode->config.walCfg.level = req.walLevel;
1486 1487 1488 1489

    walChanged = true;
  }

1490 1491
  if (pVnode->config.tsdbCfg.keep0 != req.daysToKeep0) {
    pVnode->config.tsdbCfg.keep0 = req.daysToKeep0;
1492
    if (!VND_IS_RSMA(pVnode)) {
M
Minglei Jin 已提交
1493
      tsdbChanged = true;
1494 1495 1496
    }
  }

1497 1498
  if (pVnode->config.tsdbCfg.keep1 != req.daysToKeep1) {
    pVnode->config.tsdbCfg.keep1 = req.daysToKeep1;
1499
    if (!VND_IS_RSMA(pVnode)) {
M
Minglei Jin 已提交
1500
      tsdbChanged = true;
1501 1502 1503
    }
  }

1504 1505
  if (pVnode->config.tsdbCfg.keep2 != req.daysToKeep2) {
    pVnode->config.tsdbCfg.keep2 = req.daysToKeep2;
1506
    if (!VND_IS_RSMA(pVnode)) {
M
Minglei Jin 已提交
1507
      tsdbChanged = true;
1508 1509 1510 1511
    }
  }

  if (walChanged) {
M
Minglei Jin 已提交
1512 1513 1514 1515 1516
    walAlter(pVnode->pWal, &pVnode->config.walCfg);
  }

  if (tsdbChanged) {
    tsdbSetKeepCfg(pVnode->pTsdb, &pVnode->config.tsdbCfg);
1517 1518
  }

1519 1520 1521
  return 0;
}

1522 1523 1524 1525 1526 1527
static int32_t vnodeProcessBatchDeleteReq(SVnode *pVnode, int64_t version, void *pReq, int32_t len, SRpcMsg *pRsp) {
  SBatchDeleteReq deleteReq;
  SDecoder        decoder;
  tDecoderInit(&decoder, pReq, len);
  tDecodeSBatchDeleteReq(&decoder, &deleteReq);

1528
  SMetaReader mr = {0};
1529
  metaReaderInit(&mr, pVnode->pMeta, META_READER_NOLOCK);
1530

1531 1532 1533
  int32_t sz = taosArrayGetSize(deleteReq.deleteReqs);
  for (int32_t i = 0; i < sz; i++) {
    SSingleDeleteReq *pOneReq = taosArrayGet(deleteReq.deleteReqs, i);
1534 1535
    char             *name = pOneReq->tbname;
    if (metaGetTableEntryByName(&mr, name) < 0) {
S
Shengliang Guan 已提交
1536
      vDebug("vgId:%d, stream delete msg, skip since no table: %s", pVnode->config.vgId, name);
1537 1538 1539 1540 1541
      continue;
    }

    int64_t uid = mr.me.uid;

L
Liu Jicong 已提交
1542
    int32_t code = tsdbDeleteTableData(pVnode->pTsdb, version, deleteReq.suid, uid, pOneReq->startTs, pOneReq->endTs);
1543 1544 1545
    if (code < 0) {
      terrno = code;
      vError("vgId:%d, delete error since %s, suid:%" PRId64 ", uid:%" PRId64 ", start ts:%" PRId64 ", end ts:%" PRId64,
L
Liu Jicong 已提交
1546
             TD_VID(pVnode), terrstr(), deleteReq.suid, uid, pOneReq->startTs, pOneReq->endTs);
1547
    }
1548 1549

    tDecoderClear(&mr.coder);
1550
  }
1551
  metaReaderClear(&mr);
1552
  taosArrayDestroy(deleteReq.deleteReqs);
1553 1554 1555
  return 0;
}

H
Hongze Cheng 已提交
1556 1557 1558 1559 1560
static int32_t vnodeProcessDeleteReq(SVnode *pVnode, int64_t version, void *pReq, int32_t len, SRpcMsg *pRsp) {
  int32_t     code = 0;
  SDecoder   *pCoder = &(SDecoder){0};
  SDeleteRes *pRes = &(SDeleteRes){0};

wmmhello's avatar
wmmhello 已提交
1561 1562 1563 1564 1565
  pRsp->msgType = TDMT_VND_DELETE_RSP;
  pRsp->pCont = NULL;
  pRsp->contLen = 0;
  pRsp->code = TSDB_CODE_SUCCESS;

H
Hongze Cheng 已提交
1566 1567 1568 1569 1570 1571 1572 1573
  pRes->uidList = taosArrayInit(0, sizeof(tb_uid_t));
  if (pRes->uidList == NULL) {
    code = TSDB_CODE_OUT_OF_MEMORY;
    goto _err;
  }

  tDecoderInit(pCoder, pReq, len);
  tDecodeDeleteRes(pCoder, pRes);
1574
  ASSERT(taosArrayGetSize(pRes->uidList) == 0 || (pRes->skey != 0 && pRes->ekey != 0));
H
Hongze Cheng 已提交
1575 1576 1577 1578 1579 1580 1581 1582 1583

  for (int32_t iUid = 0; iUid < taosArrayGetSize(pRes->uidList); iUid++) {
    code = tsdbDeleteTableData(pVnode->pTsdb, version, pRes->suid, *(uint64_t *)taosArrayGet(pRes->uidList, iUid),
                               pRes->skey, pRes->ekey);
    if (code) goto _err;
  }

  tDecoderClear(pCoder);
  taosArrayDestroy(pRes->uidList);
wmmhello's avatar
wmmhello 已提交
1584 1585

  SVDeleteRsp rsp = {.affectedRows = pRes->affectedRows};
L
Liu Jicong 已提交
1586
  int32_t     ret = 0;
wmmhello's avatar
wmmhello 已提交
1587 1588
  tEncodeSize(tEncodeSVDeleteRsp, &rsp, pRsp->contLen, ret);
  pRsp->pCont = rpcMallocCont(pRsp->contLen);
L
Liu Jicong 已提交
1589
  SEncoder ec = {0};
wmmhello's avatar
wmmhello 已提交
1590 1591 1592
  tEncoderInit(&ec, pRsp->pCont, pRsp->contLen);
  tEncodeSVDeleteRsp(&ec, &rsp);
  tEncoderClear(&ec);
H
Hongze Cheng 已提交
1593 1594 1595 1596
  return code;

_err:
  return code;
M
Minglei Jin 已提交
1597
}
dengyihao's avatar
dengyihao 已提交
1598
static int32_t vnodeProcessCreateIndexReq(SVnode *pVnode, int64_t version, void *pReq, int32_t len, SRpcMsg *pRsp) {
dengyihao's avatar
dengyihao 已提交
1599 1600 1601 1602 1603 1604 1605 1606 1607 1608 1609 1610 1611 1612 1613 1614 1615 1616 1617 1618 1619 1620 1621 1622
  SVCreateStbReq req = {0};
  SDecoder       dc = {0};

  pRsp->msgType = TDMT_VND_CREATE_INDEX_RSP;
  pRsp->code = TSDB_CODE_SUCCESS;
  pRsp->pCont = NULL;
  pRsp->contLen = 0;

  tDecoderInit(&dc, pReq, len);
  // decode req
  if (tDecodeSVCreateStbReq(&dc, &req) < 0) {
    terrno = TSDB_CODE_INVALID_MSG;
    tDecoderClear(&dc);
    return -1;
  }
  if (metaAddIndexToSTable(pVnode->pMeta, version, &req) < 0) {
    pRsp->code = terrno;
    goto _err;
  }
  tDecoderClear(&dc);
  return 0;
_err:
  tDecoderClear(&dc);
  return -1;
dengyihao's avatar
dengyihao 已提交
1623 1624
}
static int32_t vnodeProcessDropIndexReq(SVnode *pVnode, int64_t version, void *pReq, int32_t len, SRpcMsg *pRsp) {
dengyihao's avatar
dengyihao 已提交
1625 1626 1627 1628 1629 1630 1631 1632 1633 1634 1635 1636 1637 1638
  SDropIndexReq req = {0};
  pRsp->msgType = TDMT_VND_CREATE_INDEX_RSP;
  pRsp->code = TSDB_CODE_SUCCESS;
  pRsp->pCont = NULL;
  pRsp->contLen = 0;

  if (tDeserializeSDropIdxReq(pReq, len, &req)) {
    terrno = TSDB_CODE_INVALID_MSG;
    return -1;
  }
  if (metaDropIndexFromSTable(pVnode->pMeta, version, &req) < 0) {
    pRsp->code = terrno;
    return -1;
  }
dengyihao's avatar
dengyihao 已提交
1639 1640
  return TSDB_CODE_SUCCESS;
}
1641 1642 1643 1644 1645 1646 1647 1648 1649 1650 1651 1652 1653 1654

static int32_t vnodeProcessCompactVnodeReq(SVnode *pVnode, int64_t version, void *pReq, int32_t len, SRpcMsg *pRsp) {
  SCompactVnodeReq req = {0};
  if (tDeserializeSCompactVnodeReq(pReq, len, &req) != 0) {
    terrno = TSDB_CODE_INVALID_MSG;
    return TSDB_CODE_INVALID_MSG;
  }
  vInfo("vgId:%d, compact msg will be processed, db:%s dbUid:%" PRId64 " compactStartTime:%" PRId64, TD_VID(pVnode),
        req.db, req.dbUid, req.compactStartTime);

  vnodeAsyncCompact(pVnode);
  vnodeBegin(pVnode);

  return 0;
dengyihao's avatar
dengyihao 已提交
1655
}