vnodeSvr.c 17.0 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
#include "vnd.h"
H
Hongze Cheng 已提交
17

H
Hongze Cheng 已提交
18
static int vnodeProcessCreateStbReq(SVnode *pVnode, int64_t version, void *pReq, int len, SRpcMsg *pRsp);
H
Hongze Cheng 已提交
19
static int vnodeProcessAlterStbReq(SVnode *pVnode, void *pReq, int32_t len, SRpcMsg *pRsp);
H
Hongze Cheng 已提交
20
static int vnodeProcessDropStbReq(SVnode *pVnode, int64_t version, void *pReq, int32_t len, SRpcMsg *pRsp);
H
Hongze Cheng 已提交
21
static int vnodeProcessCreateTbReq(SVnode *pVnode, int64_t version, void *pReq, int len, SRpcMsg *pRsp);
H
Hongze Cheng 已提交
22
static int vnodeProcessAlterTbReq(SVnode *pVnode, void *pReq, int32_t len, SRpcMsg *pRsp);
H
Hongze Cheng 已提交
23
static int vnodeProcessDropTbReq(SVnode *pVnode, int64_t version, void *pReq, int32_t len, SRpcMsg *pRsp);
H
Hongze Cheng 已提交
24
static int vnodeProcessSubmitReq(SVnode *pVnode, int64_t version, void *pReq, int32_t len, SRpcMsg *pRsp);
H
Hongze Cheng 已提交
25

H
Hongze Cheng 已提交
26
int vnodePreprocessWriteReqs(SVnode *pVnode, SArray *pMsgs, int64_t *version) {
H
Hongze Cheng 已提交
27
#if 0
H
Hongze Cheng 已提交
28 29 30
  SNodeMsg *pMsg;
  SRpcMsg  *pRpc;

H
Hongze Cheng 已提交
31
  *version = pVnode->state.processed;
H
Hongze Cheng 已提交
32 33 34 35 36
  for (int i = 0; i < taosArrayGetSize(pMsgs); i++) {
    pMsg = *(SNodeMsg **)taosArrayGet(pMsgs, i);
    pRpc = &pMsg->rpcMsg;

    // set request version
H
Hongze Cheng 已提交
37
    if (walWrite(pVnode->pWal, pVnode->state.processed++, pRpc->msgType, pRpc->pCont, pRpc->contLen) < 0) {
H
refact  
Hongze Cheng 已提交
38
      vError("vnode:%d  write wal error since %s", TD_VID(pVnode), terrstr());
H
Hongze Cheng 已提交
39
      return -1;
H
Hongze Cheng 已提交
40 41 42 43
    }
  }

  walFsync(pVnode->pWal, false);
H
Hongze Cheng 已提交
44

H
Hongze Cheng 已提交
45
#endif
H
Hongze Cheng 已提交
46
  return 0;
H
Hongze Cheng 已提交
47 48
}

H
Hongze Cheng 已提交
49
int vnodeProcessWriteReq(SVnode *pVnode, SRpcMsg *pMsg, int64_t version, SRpcMsg *pRsp) {
H
Hongze Cheng 已提交
50
  void *ptr = NULL;
H
Hongze Cheng 已提交
51 52
  void *pReq;
  int   len;
H
Hongze Cheng 已提交
53
  int   ret;
H
Hongze Cheng 已提交
54

H
Hongze Cheng 已提交
55
  vTrace("vgId:%d start to process write request %s, version %" PRId64, TD_VID(pVnode), TMSG_INFO(pMsg->msgType),
H
Hongze Cheng 已提交
56
         version);
H
Hongze Cheng 已提交
57

H
Hongze Cheng 已提交
58 59
  pVnode->state.applied = version;

H
Hongze Cheng 已提交
60 61 62
  // skip header
  pReq = POINTER_SHIFT(pMsg->pCont, sizeof(SMsgHead));
  len = pMsg->contLen - sizeof(SMsgHead);
H
Hongze Cheng 已提交
63

H
Hongze Cheng 已提交
64
  if (tqPushMsg(pVnode->pTq, pMsg->pCont, pMsg->contLen, pMsg->msgType, version) < 0) {
H
Hongze Cheng 已提交
65
    vError("vgId:%d failed to push msg to TQ since %s", TD_VID(pVnode), tstrerror(terrno));
H
Hongze Cheng 已提交
66
    return -1;
H
Hongze Cheng 已提交
67 68 69
  }

  switch (pMsg->msgType) {
H
Hongze Cheng 已提交
70
    /* META */
H
Hongze Cheng 已提交
71
    case TDMT_VND_CREATE_STB:
H
Hongze Cheng 已提交
72
      if (vnodeProcessCreateStbReq(pVnode, version, pReq, len, pRsp) < 0) goto _err;
H
Hongze Cheng 已提交
73
      break;
H
Hongze Cheng 已提交
74
    case TDMT_VND_ALTER_STB:
H
Hongze Cheng 已提交
75
      if (vnodeProcessAlterStbReq(pVnode, pReq, len, pRsp) < 0) goto _err;
H
Hongze Cheng 已提交
76
      break;
H
Hongze Cheng 已提交
77
    case TDMT_VND_DROP_STB:
H
Hongze Cheng 已提交
78
      if (vnodeProcessDropStbReq(pVnode, version, pReq, len, pRsp) < 0) goto _err;
H
Hongze Cheng 已提交
79
      break;
H
Hongze Cheng 已提交
80
    case TDMT_VND_CREATE_TABLE:
H
Hongze Cheng 已提交
81
      if (vnodeProcessCreateTbReq(pVnode, version, pReq, len, pRsp) < 0) goto _err;
H
Hongze Cheng 已提交
82 83
      break;
    case TDMT_VND_ALTER_TABLE:
H
Hongze Cheng 已提交
84
      if (vnodeProcessAlterTbReq(pVnode, pReq, len, pRsp) < 0) goto _err;
H
Hongze Cheng 已提交
85
      break;
H
Hongze Cheng 已提交
86
    case TDMT_VND_DROP_TABLE:
H
Hongze Cheng 已提交
87
      if (vnodeProcessDropTbReq(pVnode, version, pReq, len, pRsp) < 0) goto _err;
H
Hongze Cheng 已提交
88
      break;
H
Hongze Cheng 已提交
89 90 91 92 93 94
    case TDMT_VND_CREATE_SMA: {  // timeRangeSMA
      if (tsdbCreateTSma(pVnode->pTsdb, POINTER_SHIFT(pMsg->pCont, sizeof(SMsgHead))) < 0) {
        // TODO
      }
    } break;
    /* TSDB */
H
Hongze Cheng 已提交
95
    case TDMT_VND_SUBMIT:
H
Hongze Cheng 已提交
96
      if (vnodeProcessSubmitReq(pVnode, version, pMsg->pCont, pMsg->contLen, pRsp) < 0) goto _err;
H
Hongze Cheng 已提交
97
      break;
H
Hongze Cheng 已提交
98
    /* TQ */
L
Liu Jicong 已提交
99 100 101 102 103 104
    case TDMT_VND_MQ_VG_CHANGE:
      if (tqProcessVgChangeReq(pVnode->pTq, POINTER_SHIFT(pMsg->pCont, sizeof(SMsgHead)),
                               pMsg->contLen - sizeof(SMsgHead)) < 0) {
        // TODO: handle error
      }
      break;
H
Hongze Cheng 已提交
105 106 107 108 109 110 111 112 113 114
    case TDMT_VND_TASK_DEPLOY: {
      if (tqProcessTaskDeploy(pVnode->pTq, POINTER_SHIFT(pMsg->pCont, sizeof(SMsgHead)),
                              pMsg->contLen - sizeof(SMsgHead)) < 0) {
      }
    } break;
    case TDMT_VND_TASK_WRITE_EXEC: {
      if (tqProcessTaskExec(pVnode->pTq, POINTER_SHIFT(pMsg->pCont, sizeof(SMsgHead)), pMsg->contLen - sizeof(SMsgHead),
                            0) < 0) {
      }
    } break;
S
Shengliang Guan 已提交
115 116
    case TDMT_VND_ALTER_VNODE:
      break;
H
Hongze Cheng 已提交
117 118 119 120 121
    default:
      ASSERT(0);
      break;
  }

H
Hongze Cheng 已提交
122
  vDebug("vgId:%d process %s request success, version: %" PRId64, TD_VID(pVnode), TMSG_INFO(pMsg->msgType), version);
H
Hongze Cheng 已提交
123

H
Hongze Cheng 已提交
124
  // commit if need
H
Hongze Cheng 已提交
125
  if (vnodeShouldCommit(pVnode)) {
H
Hongze Cheng 已提交
126
    vInfo("vgId:%d commit at version %" PRId64, TD_VID(pVnode), version);
H
Hongze Cheng 已提交
127 128 129 130 131
    // commit current change
    vnodeCommit(pVnode);

    // start a new one
    vnodeBegin(pVnode);
H
Hongze Cheng 已提交
132 133 134
  }

  return 0;
H
Hongze Cheng 已提交
135 136

_err:
H
Hongze Cheng 已提交
137
  vDebug("vgId:%d process %s request failed since %s, version: %" PRId64, TD_VID(pVnode), TMSG_INFO(pMsg->msgType),
H
Hongze Cheng 已提交
138 139
         tstrerror(terrno), version);
  return -1;
H
Hongze Cheng 已提交
140 141 142
}

int vnodeProcessQueryMsg(SVnode *pVnode, SRpcMsg *pMsg) {
D
dapan 已提交
143
  vTrace("message in vnode query queue is processing");
144
#if 0
145
  SReadHandle handle = {.reader = pVnode->pTsdb, .meta = pVnode->pMeta, .config = &pVnode->config, .vnode = pVnode};
146 147
#endif
  SReadHandle handle = {.meta = pVnode->pMeta, .config = &pVnode->config, .vnode = pVnode};
H
Hongze Cheng 已提交
148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200
  switch (pMsg->msgType) {
    case TDMT_VND_QUERY:
      return qWorkerProcessQueryMsg(&handle, pVnode->pQuery, pMsg);
    case TDMT_VND_QUERY_CONTINUE:
      return qWorkerProcessCQueryMsg(&handle, pVnode->pQuery, pMsg);
    default:
      vError("unknown msg type:%d in query queue", pMsg->msgType);
      return TSDB_CODE_VND_APP_ERROR;
  }
}

int vnodeProcessFetchMsg(SVnode *pVnode, SRpcMsg *pMsg, SQueueInfo *pInfo) {
  vTrace("message in fetch queue is processing");
  char   *msgstr = POINTER_SHIFT(pMsg->pCont, sizeof(SMsgHead));
  int32_t msgLen = pMsg->contLen - sizeof(SMsgHead);
  switch (pMsg->msgType) {
    case TDMT_VND_FETCH:
      return qWorkerProcessFetchMsg(pVnode, pVnode->pQuery, pMsg);
    case TDMT_VND_FETCH_RSP:
      return qWorkerProcessFetchRsp(pVnode, pVnode->pQuery, pMsg);
    case TDMT_VND_RES_READY:
      return qWorkerProcessReadyMsg(pVnode, pVnode->pQuery, pMsg);
    case TDMT_VND_TASKS_STATUS:
      return qWorkerProcessStatusMsg(pVnode, pVnode->pQuery, pMsg);
    case TDMT_VND_CANCEL_TASK:
      return qWorkerProcessCancelMsg(pVnode, pVnode->pQuery, pMsg);
    case TDMT_VND_DROP_TASK:
      return qWorkerProcessDropMsg(pVnode, pVnode->pQuery, pMsg);
    case TDMT_VND_TABLE_META:
      return vnodeGetTableMeta(pVnode, pMsg);
    case TDMT_VND_CONSUME:
      return tqProcessPollReq(pVnode->pTq, pMsg, pInfo->workerId);
    case TDMT_VND_TASK_PIPE_EXEC:
    case TDMT_VND_TASK_MERGE_EXEC:
      return tqProcessTaskExec(pVnode->pTq, msgstr, msgLen, 0);
    case TDMT_VND_STREAM_TRIGGER:
      return tqProcessStreamTrigger(pVnode->pTq, pMsg->pCont, pMsg->contLen, 0);
    case TDMT_VND_QUERY_HEARTBEAT:
      return qWorkerProcessHbMsg(pVnode, pVnode->pQuery, pMsg);
    default:
      vError("unknown msg type:%d in fetch queue", pMsg->msgType);
      return TSDB_CODE_VND_APP_ERROR;
  }
}

// TODO: remove the function
void smaHandleRes(void *pVnode, int64_t smaId, const SArray *data) {
  // TODO

  // blockDebugShowData(data);
  tsdbInsertTSmaData(((SVnode *)pVnode)->pTsdb, smaId, (const char *)data);
}

M
Minghao Li 已提交
201
// sync integration
H
Hongze Cheng 已提交
202
int vnodeProcessSyncReq(SVnode *pVnode, SRpcMsg *pMsg, SRpcMsg **pRsp) {
203 204 205
  if (syncEnvIsStart()) {
    SSyncNode *pSyncNode = syncNodeAcquire(pVnode->sync);
    assert(pSyncNode != NULL);
M
Minghao Li 已提交
206

207 208
    ESyncState state = syncGetMyRole(pVnode->sync);
    SyncTerm   currentTerm = syncGetMyTerm(pVnode->sync);
M
Minghao Li 已提交
209

210
    SMsgHead *pHead = pMsg->pCont;
M
Minghao Li 已提交
211

212 213 214 215 216
    char  logBuf[512];
    char *syncNodeStr = sync2SimpleStr(pVnode->sync);
    snprintf(logBuf, sizeof(logBuf), "==vnodeProcessSyncReq== msgType:%d, syncNode: %s", pMsg->msgType, syncNodeStr);
    syncRpcMsgLog2(logBuf, pMsg);
    taosMemoryFree(syncNodeStr);
M
Minghao Li 已提交
217

218
    SRpcMsg *pRpcMsg = pMsg;
M
Minghao Li 已提交
219

220 221 222
    if (pRpcMsg->msgType == TDMT_VND_SYNC_TIMEOUT) {
      SyncTimeout *pSyncMsg = syncTimeoutFromRpcMsg2(pRpcMsg);
      assert(pSyncMsg != NULL);
M
Minghao Li 已提交
223

224 225
      syncNodeOnTimeoutCb(pSyncNode, pSyncMsg);
      syncTimeoutDestroy(pSyncMsg);
M
Minghao Li 已提交
226

227 228 229
    } else if (pRpcMsg->msgType == TDMT_VND_SYNC_PING) {
      SyncPing *pSyncMsg = syncPingFromRpcMsg2(pRpcMsg);
      assert(pSyncMsg != NULL);
M
Minghao Li 已提交
230

231 232
      syncNodeOnPingCb(pSyncNode, pSyncMsg);
      syncPingDestroy(pSyncMsg);
M
Minghao Li 已提交
233

234 235 236
    } else if (pRpcMsg->msgType == TDMT_VND_SYNC_PING_REPLY) {
      SyncPingReply *pSyncMsg = syncPingReplyFromRpcMsg2(pRpcMsg);
      assert(pSyncMsg != NULL);
M
Minghao Li 已提交
237

238 239
      syncNodeOnPingReplyCb(pSyncNode, pSyncMsg);
      syncPingReplyDestroy(pSyncMsg);
M
Minghao Li 已提交
240

241 242 243
    } else if (pRpcMsg->msgType == TDMT_VND_SYNC_CLIENT_REQUEST) {
      SyncClientRequest *pSyncMsg = syncClientRequestFromRpcMsg2(pRpcMsg);
      assert(pSyncMsg != NULL);
M
Minghao Li 已提交
244

245 246
      syncNodeOnClientRequestCb(pSyncNode, pSyncMsg);
      syncClientRequestDestroy(pSyncMsg);
M
Minghao Li 已提交
247

248 249 250
    } else if (pRpcMsg->msgType == TDMT_VND_SYNC_REQUEST_VOTE) {
      SyncRequestVote *pSyncMsg = syncRequestVoteFromRpcMsg2(pRpcMsg);
      assert(pSyncMsg != NULL);
M
Minghao Li 已提交
251

252 253
      syncNodeOnRequestVoteCb(pSyncNode, pSyncMsg);
      syncRequestVoteDestroy(pSyncMsg);
M
Minghao Li 已提交
254

255 256 257
    } else if (pRpcMsg->msgType == TDMT_VND_SYNC_REQUEST_VOTE_REPLY) {
      SyncRequestVoteReply *pSyncMsg = syncRequestVoteReplyFromRpcMsg2(pRpcMsg);
      assert(pSyncMsg != NULL);
M
Minghao Li 已提交
258

259 260
      syncNodeOnRequestVoteReplyCb(pSyncNode, pSyncMsg);
      syncRequestVoteReplyDestroy(pSyncMsg);
M
Minghao Li 已提交
261

262 263 264
    } else if (pRpcMsg->msgType == TDMT_VND_SYNC_APPEND_ENTRIES) {
      SyncAppendEntries *pSyncMsg = syncAppendEntriesFromRpcMsg2(pRpcMsg);
      assert(pSyncMsg != NULL);
M
Minghao Li 已提交
265

266 267
      syncNodeOnAppendEntriesCb(pSyncNode, pSyncMsg);
      syncAppendEntriesDestroy(pSyncMsg);
M
Minghao Li 已提交
268

269 270 271 272 273 274
    } else if (pRpcMsg->msgType == TDMT_VND_SYNC_APPEND_ENTRIES_REPLY) {
      SyncAppendEntriesReply *pSyncMsg = syncAppendEntriesReplyFromRpcMsg2(pRpcMsg);
      assert(pSyncMsg != NULL);

      syncNodeOnAppendEntriesReplyCb(pSyncNode, pSyncMsg);
      syncAppendEntriesReplyDestroy(pSyncMsg);
M
Minghao Li 已提交
275

276 277 278
    } else {
      vError("==vnodeProcessSyncReq== error msg type:%d", pRpcMsg->msgType);
    }
M
Minghao Li 已提交
279

280 281 282 283
    syncNodeRelease(pSyncNode);
  } else {
    vError("==vnodeProcessSyncReq== error syncEnv stop");
  }
H
Hongze Cheng 已提交
284
  return 0;
H
Hongze Cheng 已提交
285 286
}

H
Hongze Cheng 已提交
287
static int vnodeProcessCreateStbReq(SVnode *pVnode, int64_t version, void *pReq, int len, SRpcMsg *pRsp) {
H
Hongze Cheng 已提交
288 289 290
  SVCreateStbReq req = {0};
  SCoder         coder;

H
Hongze Cheng 已提交
291 292 293 294 295 296
  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 已提交
297 298 299
  tCoderInit(&coder, TD_LITTLE_ENDIAN, pReq, len, TD_DECODER);

  if (tDecodeSVCreateStbReq(&coder, &req) < 0) {
H
Hongze Cheng 已提交
300 301
    pRsp->code = terrno;
    goto _err;
H
Hongze Cheng 已提交
302 303
  }

H
Hongze Cheng 已提交
304
  if (metaCreateSTable(pVnode->pMeta, version, &req) < 0) {
H
Hongze Cheng 已提交
305 306
    pRsp->code = terrno;
    goto _err;
H
Hongze Cheng 已提交
307 308
  }

C
Cary Xu 已提交
309 310
  tsdbRegisterRSma(pVnode->pTsdb, pVnode->pMeta, &req);

H
Hongze Cheng 已提交
311
  tCoderClear(&coder);
H
Hongze Cheng 已提交
312
  return 0;
H
Hongze Cheng 已提交
313 314 315 316

_err:
  tCoderClear(&coder);
  return -1;
H
Hongze Cheng 已提交
317 318
}

H
Hongze Cheng 已提交
319 320 321 322 323
static int vnodeProcessCreateTbReq(SVnode *pVnode, int64_t version, void *pReq, int len, SRpcMsg *pRsp) {
  SCoder             coder = {0};
  int                rcode = 0;
  SVCreateTbBatchReq req = {0};
  SVCreateTbReq     *pCreateReq;
H
Hongze Cheng 已提交
324 325 326
  SVCreateTbBatchRsp rsp = {0};
  SVCreateTbRsp      cRsp = {0};
  char               tbName[TSDB_TABLE_FNAME_LEN];
C
Cary Xu 已提交
327
  STbUidStore       *pStore = NULL;
H
Hongze Cheng 已提交
328 329

  pRsp->msgType = TDMT_VND_CREATE_TABLE_RSP;
H
Hongze Cheng 已提交
330 331 332
  pRsp->code = TSDB_CODE_SUCCESS;
  pRsp->pCont = NULL;
  pRsp->contLen = 0;
H
Hongze Cheng 已提交
333

H
Hongze Cheng 已提交
334 335 336 337 338 339 340
  // decode
  tCoderInit(&coder, TD_LITTLE_ENDIAN, pReq, len, TD_DECODER);
  if (tDecodeSVCreateTbBatchReq(&coder, &req) < 0) {
    rcode = -1;
    terrno = TSDB_CODE_INVALID_MSG;
    goto _exit;
  }
H
Hongze Cheng 已提交
341

H
Hongze Cheng 已提交
342 343 344 345 346 347 348
  rsp.pArray = taosArrayInit(sizeof(cRsp), req.nReqs);
  if (rsp.pArray == NULL) {
    rcode = -1;
    terrno = TSDB_CODE_OUT_OF_MEMORY;
    goto _exit;
  }

H
Hongze Cheng 已提交
349 350 351
  // loop to create table
  for (int iReq = 0; iReq < req.nReqs; iReq++) {
    pCreateReq = req.pReqs + iReq;
H
Hongze Cheng 已提交
352 353 354 355 356 357 358 359 360 361

    // 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
H
Hongze Cheng 已提交
362
    if (metaCreateTable(pVnode->pMeta, version, pCreateReq) < 0) {
H
Hongze Cheng 已提交
363
      cRsp.code = terrno;
H
Hongze Cheng 已提交
364
    } else {
H
Hongze Cheng 已提交
365
      cRsp.code = TSDB_CODE_SUCCESS;
C
Cary Xu 已提交
366
      tsdbFetchTbUidList(pVnode->pTsdb, &pStore, pCreateReq->ctb.suid, pCreateReq->uid);
H
Hongze Cheng 已提交
367
    }
H
Hongze Cheng 已提交
368 369

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

H
Hongze Cheng 已提交
372 373
  tCoderClear(&coder);

C
Cary Xu 已提交
374 375 376
  tsdbUpdateTbUidList(pVnode->pTsdb, pStore);
  tsdbUidStoreFree(pStore);

H
Hongze Cheng 已提交
377
  // prepare rsp
wafwerar's avatar
wafwerar 已提交
378 379
  int32_t ret = 0;
  tEncodeSize(tEncodeSVCreateTbBatchRsp, &rsp, pRsp->contLen, ret);
H
Hongze Cheng 已提交
380 381 382 383 384 385 386 387
  pRsp->pCont = rpcMallocCont(pRsp->contLen);
  if (pRsp->pCont == NULL) {
    terrno = TSDB_CODE_OUT_OF_MEMORY;
    rcode = -1;
    goto _exit;
  }
  tCoderInit(&coder, TD_LITTLE_ENDIAN, pRsp->pCont, pRsp->contLen, TD_ENCODER);
  tEncodeSVCreateTbBatchRsp(&coder, &rsp);
H
Hongze Cheng 已提交
388

H
Hongze Cheng 已提交
389
_exit:
H
Hongze Cheng 已提交
390
  taosArrayClear(rsp.pArray);
H
Hongze Cheng 已提交
391 392
  tCoderClear(&coder);
  return rcode;
H
Hongze Cheng 已提交
393 394
}

H
Hongze Cheng 已提交
395
static int vnodeProcessAlterStbReq(SVnode *pVnode, void *pReq, int32_t len, SRpcMsg *pRsp) {
H
Hongze Cheng 已提交
396
  // ASSERT(0);
H
Hongze Cheng 已提交
397
#if 0
H
Hongze Cheng 已提交
398
  SVCreateTbReq vAlterTbReq = {0};
H
refact  
Hongze Cheng 已提交
399
  vTrace("vgId:%d, process alter stb req", TD_VID(pVnode));
H
Hongze Cheng 已提交
400 401 402 403 404 405 406 407
  tDeserializeSVCreateTbReq(pReq, &vAlterTbReq);
  // TODO: to encapsule a free API
  taosMemoryFree(vAlterTbReq.stbCfg.pSchema);
  taosMemoryFree(vAlterTbReq.stbCfg.pTagSchema);
  if (vAlterTbReq.stbCfg.pRSmaParam) {
    taosMemoryFree(vAlterTbReq.stbCfg.pRSmaParam);
  }
  taosMemoryFree(vAlterTbReq.name);
H
Hongze Cheng 已提交
408 409 410 411
#endif
  return 0;
}

H
Hongze Cheng 已提交
412
static int vnodeProcessDropStbReq(SVnode *pVnode, int64_t version, void *pReq, int32_t len, SRpcMsg *pRsp) {
H
Hongze Cheng 已提交
413 414 415 416 417 418 419 420 421 422 423 424 425 426 427 428
  SVDropStbReq req = {0};
  int          rcode = TSDB_CODE_SUCCESS;
  SCoder       coder = {0};

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

  // decode request
  tCoderInit(&coder, TD_LITTLE_ENDIAN, pReq, len, TD_DECODER);
  if (tDecodeSVDropStbReq(&coder, &req) < 0) {
    rcode = TSDB_CODE_INVALID_MSG;
    goto _exit;
  }

  // process request
H
Hongze Cheng 已提交
429 430 431 432
  // if (metaDropSTable(pVnode->pMeta, version, &req) < 0) {
  //   rcode = terrno;
  //   goto _exit;
  // }
H
Hongze Cheng 已提交
433 434 435 436 437

  // return rsp
_exit:
  pRsp->code = rcode;
  tCoderClear(&coder);
H
Hongze Cheng 已提交
438 439 440 441 442 443 444 445 446
  return 0;
}

static int vnodeProcessAlterTbReq(SVnode *pVnode, void *pReq, int32_t len, SRpcMsg *pRsp) {
  // TODO
  ASSERT(0);
  return 0;
}

H
Hongze Cheng 已提交
447
static int vnodeProcessDropTbReq(SVnode *pVnode, int64_t version, void *pReq, int32_t len, SRpcMsg *pRsp) {
H
Hongze Cheng 已提交
448 449 450 451 452
  SVDropTbBatchReq req = {0};
  SVDropTbBatchRsp rsp = {0};
  SCoder           coder = {0};
  int              ret;

H
Hongze Cheng 已提交
453
  pRsp->msgType = TDMT_VND_DROP_TABLE_RSP;
H
Hongze Cheng 已提交
454 455 456
  pRsp->pCont = NULL;
  pRsp->contLen = 0;
  pRsp->code = TSDB_CODE_SUCCESS;
H
Hongze Cheng 已提交
457 458

  // decode req
H
Hongze Cheng 已提交
459 460 461 462 463 464 465
  tCoderInit(&coder, TD_LITTLE_ENDIAN, pReq, len, TD_DECODER);
  ret = tDecodeSVDropTbBatchReq(&coder, &req);
  if (ret < 0) {
    terrno = TSDB_CODE_INVALID_MSG;
    pRsp->code = terrno;
    goto _exit;
  }
H
Hongze Cheng 已提交
466 467

  // process req
H
Hongze Cheng 已提交
468 469 470 471
  rsp.pArray = taosArrayInit(sizeof(SVDropTbRsp), req.nReqs);
  for (int iReq = 0; iReq < req.nReqs; iReq++) {
    SVDropTbReq *pDropTbReq = req.pReqs + iReq;
    SVDropTbRsp  dropTbRsp = {0};
H
Hongze Cheng 已提交
472

H
Hongze Cheng 已提交
473 474 475
    /* code */
    ret = metaDropTable(pVnode->pMeta, version, pDropTbReq);
    if (ret < 0) {
H
Hongze Cheng 已提交
476 477 478 479 480
      if (pDropTbReq->igNotExists && terrno == TSDB_CODE_VND_TABLE_NOT_EXIST) {
        dropTbRsp.code = TSDB_CODE_SUCCESS;
      } else {
        dropTbRsp.code = terrno;
      }
H
Hongze Cheng 已提交
481
    } else {
H
Hongze Cheng 已提交
482
      dropTbRsp.code = TSDB_CODE_SUCCESS;
H
Hongze Cheng 已提交
483 484 485 486 487 488 489 490
    }

    taosArrayPush(rsp.pArray, &dropTbRsp);
  }

_exit:
  tCoderClear(&coder);
  // encode rsp (TODO)
H
Hongze Cheng 已提交
491 492 493
  return 0;
}

H
Hongze Cheng 已提交
494
static int vnodeProcessSubmitReq(SVnode *pVnode, int64_t version, void *pReq, int32_t len, SRpcMsg *pRsp) {
H
Hongze Cheng 已提交
495 496 497 498 499 500 501
  SSubmitReq    *pSubmitReq = (SSubmitReq *)pReq;
  SSubmitMsgIter msgIter = {0};
  SSubmitBlk    *pBlock;
  SSubmitRsp     rsp = {0};
  SVCreateTbReq  createTbReq = {0};
  SCoder         coder = {0};
  int32_t        nRows;
H
Hongze Cheng 已提交
502 503

  pRsp->code = 0;
C
Cary Xu 已提交
504

H
Hongze Cheng 已提交
505
  // handle the request
H
Hongze Cheng 已提交
506 507 508
  if (tInitSubmitMsgIter(pSubmitReq, &msgIter) < 0) {
    pRsp->code = TSDB_CODE_INVALID_MSG;
    goto _exit;
H
Hongze Cheng 已提交
509 510
  }

H
Hongze Cheng 已提交
511 512 513 514 515 516 517 518 519 520 521 522 523 524 525 526 527 528 529 530 531 532 533 534 535 536 537 538 539 540 541
  for (;;) {
    tGetSubmitMsgNext(&msgIter, &pBlock);
    if (pBlock == NULL) break;

    // create table for auto create table mode
    if (msgIter.schemaLen > 0) {
      tCoderInit(&coder, TD_LITTLE_ENDIAN, pBlock->data, msgIter.schemaLen, TD_DECODER);
      if (tDecodeSVCreateTbReq(&coder, &createTbReq) < 0) {
        pRsp->code = TSDB_CODE_INVALID_MSG;
        tCoderClear(&coder);
        goto _exit;
      }

      if (metaCreateTable(pVnode->pMeta, version, &createTbReq) < 0) {
        if (terrno != TSDB_CODE_TDB_TABLE_ALREADY_EXIST) {
          pRsp->code = terrno;
          tCoderClear(&coder);
          goto _exit;
        }
      }

      tCoderClear(&coder);
    }

    if (tsdbInsertTableData(pVnode->pTsdb, &msgIter, pBlock, &nRows) < 0) {
      pRsp->code = terrno;
      goto _exit;
    }

    rsp.numOfRows += nRows;
  }
542

H
Hongze Cheng 已提交
543
_exit:
H
Hongze Cheng 已提交
544 545 546 547 548
  // encode the response (TODO)
  pRsp->pCont = rpcMallocCont(sizeof(SSubmitRsp));
  memcpy(pRsp->pCont, &rsp, sizeof(rsp));
  pRsp->contLen = sizeof(SSubmitRsp);

C
Cary Xu 已提交
549 550
  tsdbTriggerRSma(pVnode->pTsdb, pReq, STREAM_DATA_TYPE_SUBMIT_BLOCK);

H
Hongze Cheng 已提交
551
  return 0;
L
Liu Jicong 已提交
552
}