vnodeSvr.c 18.6 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 201
  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);
}

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

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

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

211 212 213 214 215
    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 已提交
216

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

268 269 270 271 272 273
    } 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 已提交
274

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

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

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

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

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

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

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

H
Hongze Cheng 已提交
310
  tDecoderClear(&coder);
H
Hongze Cheng 已提交
311
  return 0;
H
Hongze Cheng 已提交
312 313

_err:
H
Hongze Cheng 已提交
314
  tDecoderClear(&coder);
H
Hongze Cheng 已提交
315
  return -1;
H
Hongze Cheng 已提交
316 317
}

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

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

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

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

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

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

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

H
Hongze Cheng 已提交
375
  tDecoderClear(&decoder);
H
Hongze Cheng 已提交
376

C
Cary Xu 已提交
377 378 379
  tsdbUpdateTbUidList(pVnode->pTsdb, pStore);
  tsdbUidStoreFree(pStore);

H
Hongze Cheng 已提交
380
  // prepare rsp
H
Hongze Cheng 已提交
381 382
  SEncoder encoder = {0};
  int32_t  ret = 0;
wafwerar's avatar
wafwerar 已提交
383
  tEncodeSize(tEncodeSVCreateTbBatchRsp, &rsp, pRsp->contLen, ret);
H
Hongze Cheng 已提交
384 385 386 387 388 389
  pRsp->pCont = rpcMallocCont(pRsp->contLen);
  if (pRsp->pCont == NULL) {
    terrno = TSDB_CODE_OUT_OF_MEMORY;
    rcode = -1;
    goto _exit;
  }
H
Hongze Cheng 已提交
390 391 392
  tEncoderInit(&encoder, pRsp->pCont, pRsp->contLen);
  tEncodeSVCreateTbBatchRsp(&encoder, &rsp);
  tEncoderClear(&encoder);
H
Hongze Cheng 已提交
393

H
Hongze Cheng 已提交
394
_exit:
H
Hongze Cheng 已提交
395
  taosArrayClear(rsp.pArray);
H
Hongze Cheng 已提交
396 397
  tDecoderClear(&decoder);
  tEncoderClear(&encoder);
H
Hongze Cheng 已提交
398
  return rcode;
H
Hongze Cheng 已提交
399 400
}

H
Hongze Cheng 已提交
401
static int vnodeProcessAlterStbReq(SVnode *pVnode, void *pReq, int32_t len, SRpcMsg *pRsp) {
H
Hongze Cheng 已提交
402
  // ASSERT(0);
H
Hongze Cheng 已提交
403
#if 0
H
Hongze Cheng 已提交
404
  SVCreateTbReq vAlterTbReq = {0};
H
refact  
Hongze Cheng 已提交
405
  vTrace("vgId:%d, process alter stb req", TD_VID(pVnode));
H
Hongze Cheng 已提交
406 407 408 409 410 411 412 413
  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 已提交
414 415 416 417
#endif
  return 0;
}

H
Hongze Cheng 已提交
418
static int vnodeProcessDropStbReq(SVnode *pVnode, int64_t version, void *pReq, int32_t len, SRpcMsg *pRsp) {
H
Hongze Cheng 已提交
419 420
  SVDropStbReq req = {0};
  int          rcode = TSDB_CODE_SUCCESS;
H
Hongze Cheng 已提交
421
  SDecoder     decoder = {0};
H
Hongze Cheng 已提交
422 423 424 425 426 427

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

  // decode request
H
Hongze Cheng 已提交
428 429
  tDecoderInit(&decoder, pReq, len);
  if (tDecodeSVDropStbReq(&decoder, &req) < 0) {
H
Hongze Cheng 已提交
430 431 432 433 434
    rcode = TSDB_CODE_INVALID_MSG;
    goto _exit;
  }

  // process request
H
Hongze Cheng 已提交
435 436 437 438
  // if (metaDropSTable(pVnode->pMeta, version, &req) < 0) {
  //   rcode = terrno;
  //   goto _exit;
  // }
H
Hongze Cheng 已提交
439 440 441 442

  // return rsp
_exit:
  pRsp->code = rcode;
H
Hongze Cheng 已提交
443
  tDecoderClear(&decoder);
H
Hongze Cheng 已提交
444 445 446 447 448 449 450 451 452
  return 0;
}

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

H
Hongze Cheng 已提交
453
static int vnodeProcessDropTbReq(SVnode *pVnode, int64_t version, void *pReq, int32_t len, SRpcMsg *pRsp) {
H
Hongze Cheng 已提交
454 455
  SVDropTbBatchReq req = {0};
  SVDropTbBatchRsp rsp = {0};
H
Hongze Cheng 已提交
456
  SDecoder         decoder = {0};
H
Hongze Cheng 已提交
457 458
  int              ret;

H
Hongze Cheng 已提交
459
  pRsp->msgType = TDMT_VND_DROP_TABLE_RSP;
H
Hongze Cheng 已提交
460 461 462
  pRsp->pCont = NULL;
  pRsp->contLen = 0;
  pRsp->code = TSDB_CODE_SUCCESS;
H
Hongze Cheng 已提交
463 464

  // decode req
H
Hongze Cheng 已提交
465 466
  tDecoderInit(&decoder, pReq, len);
  ret = tDecodeSVDropTbBatchReq(&decoder, &req);
H
Hongze Cheng 已提交
467 468 469 470 471
  if (ret < 0) {
    terrno = TSDB_CODE_INVALID_MSG;
    pRsp->code = terrno;
    goto _exit;
  }
H
Hongze Cheng 已提交
472 473

  // process req
H
Hongze Cheng 已提交
474 475 476 477
  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 已提交
478

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

    taosArrayPush(rsp.pArray, &dropTbRsp);
  }

_exit:
H
Hongze Cheng 已提交
495
  tDecoderClear(&decoder);
H
Hongze Cheng 已提交
496
  // encode rsp (TODO)
H
Hongze Cheng 已提交
497 498 499
  return 0;
}

C
Cary Xu 已提交
500 501 502 503 504 505 506 507 508 509 510 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
static int vnodeDebugPrintSubmitMsg(SVnode *pVnode, SSubmitReq *pMsg, const char* tags) {
  ASSERT(pMsg != NULL);
  SSubmitMsgIter msgIter = {0};
  SMeta         *pMeta = pVnode->pMeta;
  SSubmitBlk    *pBlock = NULL;
  SSubmitBlkIter blkIter = {0};
  STSRow        *row = NULL;
  STSchema      *pSchema = NULL;
  tb_uid_t       suid = 0;

  if (tInitSubmitMsgIter(pMsg, &msgIter) < 0) return -1;
  while (true) {
    if (tGetSubmitMsgNext(&msgIter, &pBlock) < 0) return -1;
    if (pBlock == NULL) break;
    tInitSubmitBlkIter(&msgIter, pBlock, &blkIter);
    if (blkIter.row == NULL) continue;
    if (!pSchema || (suid != msgIter.suid)) {
      if (pSchema) {
        taosMemoryFreeClear(pSchema);
      }
      pSchema = metaGetTbTSchema(pMeta, msgIter.suid, 0);  // TODO: use the real schema
      if(pSchema) {
        suid = msgIter.suid;
      }
    }
    if(!pSchema) {
      printf("%s:%d no valid schema\n", tags, __LINE__);
      continue;
    }
    char __tags[128] = {0};
    snprintf(__tags, 128, "%s: uid %" PRIi64 " ", tags, msgIter.uid);
    while ((row = tGetSubmitBlkNext(&blkIter))) {
      tdSRowPrint(row, pSchema, __tags);
    }
  }

  taosMemoryFreeClear(pSchema);

  return 0;
}

H
Hongze Cheng 已提交
541
static int vnodeProcessSubmitReq(SVnode *pVnode, int64_t version, void *pReq, int32_t len, SRpcMsg *pRsp) {
H
Hongze Cheng 已提交
542 543 544 545 546
  SSubmitReq    *pSubmitReq = (SSubmitReq *)pReq;
  SSubmitMsgIter msgIter = {0};
  SSubmitBlk    *pBlock;
  SSubmitRsp     rsp = {0};
  SVCreateTbReq  createTbReq = {0};
H
Hongze Cheng 已提交
547
  SDecoder       decoder = {0};
H
Hongze Cheng 已提交
548
  int32_t        nRows;
H
Hongze Cheng 已提交
549 550

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

C
Cary Xu 已提交
552 553 554 555
#ifdef TD_DEBUG_PRINT_ROW
  vnodeDebugPrintSubmitMsg(pVnode, pReq, __func__);
#endif

H
Hongze Cheng 已提交
556
  // handle the request
H
Hongze Cheng 已提交
557 558 559
  if (tInitSubmitMsgIter(pSubmitReq, &msgIter) < 0) {
    pRsp->code = TSDB_CODE_INVALID_MSG;
    goto _exit;
H
Hongze Cheng 已提交
560 561
  }

H
Hongze Cheng 已提交
562 563 564 565 566 567
  for (;;) {
    tGetSubmitMsgNext(&msgIter, &pBlock);
    if (pBlock == NULL) break;

    // create table for auto create table mode
    if (msgIter.schemaLen > 0) {
H
Hongze Cheng 已提交
568 569
      tDecoderInit(&decoder, pBlock->data, msgIter.schemaLen);
      if (tDecodeSVCreateTbReq(&decoder, &createTbReq) < 0) {
H
Hongze Cheng 已提交
570
        pRsp->code = TSDB_CODE_INVALID_MSG;
H
Hongze Cheng 已提交
571
        tDecoderClear(&decoder);
H
Hongze Cheng 已提交
572 573 574 575 576 577
        goto _exit;
      }

      if (metaCreateTable(pVnode->pMeta, version, &createTbReq) < 0) {
        if (terrno != TSDB_CODE_TDB_TABLE_ALREADY_EXIST) {
          pRsp->code = terrno;
H
Hongze Cheng 已提交
578
          tDecoderClear(&decoder);
H
Hongze Cheng 已提交
579 580 581 582
          goto _exit;
        }
      }

H
Hongze Cheng 已提交
583 584 585 586 587 588 589
      msgIter.uid = createTbReq.uid;
      if (createTbReq.type == TSDB_CHILD_TABLE) {
        msgIter.suid = createTbReq.ctb.suid;
      } else {
        msgIter.suid = 0;
      }

H
Hongze Cheng 已提交
590
      tDecoderClear(&decoder);
H
Hongze Cheng 已提交
591 592 593 594 595 596 597
    }

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

D
fix bug  
dapan1121 已提交
598
    rsp.affectedRows += nRows;
C
Cary Xu 已提交
599
    
H
Hongze Cheng 已提交
600
  }
601

H
Hongze Cheng 已提交
602
_exit:
H
Hongze Cheng 已提交
603 604 605 606 607
  // encode the response (TODO)
  pRsp->pCont = rpcMallocCont(sizeof(SSubmitRsp));
  memcpy(pRsp->pCont, &rsp, sizeof(rsp));
  pRsp->contLen = sizeof(SSubmitRsp);

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

H
Hongze Cheng 已提交
610
  return 0;
L
Liu Jicong 已提交
611
}