mndStream.c 33.8 KB
Newer Older
L
Liu Jicong 已提交
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19
/*
 * 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/>.
 */

#include "mndStream.h"
#include "mndDb.h"
#include "mndDnode.h"
#include "mndMnode.h"
L
Liu Jicong 已提交
20
#include "mndPrivilege.h"
L
Liu Jicong 已提交
21
#include "mndScheduler.h"
L
Liu Jicong 已提交
22 23
#include "mndShow.h"
#include "mndStb.h"
L
Liu Jicong 已提交
24
#include "mndTopic.h"
L
Liu Jicong 已提交
25 26 27
#include "mndTrans.h"
#include "mndUser.h"
#include "mndVgroup.h"
L
Liu Jicong 已提交
28
#include "parser.h"
L
Liu Jicong 已提交
29 30 31 32 33 34 35 36
#include "tname.h"

#define MND_STREAM_VER_NUMBER   1
#define MND_STREAM_RESERVE_SIZE 64

static int32_t mndStreamActionInsert(SSdb *pSdb, SStreamObj *pStream);
static int32_t mndStreamActionDelete(SSdb *pSdb, SStreamObj *pStream);
static int32_t mndStreamActionUpdate(SSdb *pSdb, SStreamObj *pStream, SStreamObj *pNewStream);
S
Shengliang Guan 已提交
37
static int32_t mndProcessCreateStreamReq(SRpcMsg *pReq);
L
Liu Jicong 已提交
38
static int32_t mndProcessDropStreamReq(SRpcMsg *pReq);
L
Liu Jicong 已提交
39
/*static int32_t mndProcessRecoverStreamReq(SRpcMsg *pReq);*/
S
Shengliang Guan 已提交
40 41 42
static int32_t mndProcessStreamMetaReq(SRpcMsg *pReq);
static int32_t mndGetStreamMeta(SRpcMsg *pReq, SShowObj *pShow, STableMetaRsp *pMeta);
static int32_t mndRetrieveStream(SRpcMsg *pReq, SShowObj *pShow, SSDataBlock *pBlock, int32_t rows);
L
Liu Jicong 已提交
43
static void    mndCancelGetNextStream(SMnode *pMnode, void *pIter);
44 45
static int32_t mndRetrieveStreamTask(SRpcMsg *pReq, SShowObj *pShow, SSDataBlock *pBlock, int32_t rows);
static void    mndCancelGetNextStreamTask(SMnode *pMnode, void *pIter);
L
Liu Jicong 已提交
46 47

int32_t mndInitStream(SMnode *pMnode) {
48 49 50 51 52 53 54 55 56
  SSdbTable table = {
      .sdbType = SDB_STREAM,
      .keyType = SDB_KEY_BINARY,
      .encodeFp = (SdbEncodeFp)mndStreamActionEncode,
      .decodeFp = (SdbDecodeFp)mndStreamActionDecode,
      .insertFp = (SdbInsertFp)mndStreamActionInsert,
      .updateFp = (SdbUpdateFp)mndStreamActionUpdate,
      .deleteFp = (SdbDeleteFp)mndStreamActionDelete,
  };
L
Liu Jicong 已提交
57 58

  mndSetMsgHandle(pMnode, TDMT_MND_CREATE_STREAM, mndProcessCreateStreamReq);
L
Liu Jicong 已提交
59
  mndSetMsgHandle(pMnode, TDMT_MND_DROP_STREAM, mndProcessDropStreamReq);
L
Liu Jicong 已提交
60
  /*mndSetMsgHandle(pMnode, TDMT_MND_RECOVER_STREAM, mndProcessRecoverStreamReq);*/
L
Liu Jicong 已提交
61 62 63

  mndSetMsgHandle(pMnode, TDMT_STREAM_TASK_DEPLOY_RSP, mndTransProcessRsp);
  mndSetMsgHandle(pMnode, TDMT_STREAM_TASK_DROP_RSP, mndTransProcessRsp);
L
Liu Jicong 已提交
64

L
Liu Jicong 已提交
65 66
  mndAddShowRetrieveHandle(pMnode, TSDB_MGMT_TABLE_STREAMS, mndRetrieveStream);
  mndAddShowFreeIterHandle(pMnode, TSDB_MGMT_TABLE_STREAMS, mndCancelGetNextStream);
67 68
  mndAddShowRetrieveHandle(pMnode, TSDB_MGMT_TABLE_STREAM_TASKS, mndRetrieveStreamTask);
  mndAddShowFreeIterHandle(pMnode, TSDB_MGMT_TABLE_STREAM_TASKS, mndCancelGetNextStreamTask);
L
Liu Jicong 已提交
69 70 71 72 73 74 75 76 77 78

  return sdbSetTable(pMnode->pSdb, table);
}

void mndCleanupStream(SMnode *pMnode) {}

SSdbRaw *mndStreamActionEncode(SStreamObj *pStream) {
  terrno = TSDB_CODE_OUT_OF_MEMORY;
  void *buf = NULL;

H
Hongze Cheng 已提交
79 80
  SEncoder encoder;
  tEncoderInit(&encoder, NULL, 0);
L
Liu Jicong 已提交
81
  if (tEncodeSStreamObj(&encoder, pStream) < 0) {
H
Hongze Cheng 已提交
82
    tEncoderClear(&encoder);
L
Liu Jicong 已提交
83 84 85
    goto STREAM_ENCODE_OVER;
  }
  int32_t tlen = encoder.pos;
H
Hongze Cheng 已提交
86
  tEncoderClear(&encoder);
L
Liu Jicong 已提交
87 88 89 90 91

  int32_t  size = sizeof(int32_t) + tlen + MND_STREAM_RESERVE_SIZE;
  SSdbRaw *pRaw = sdbAllocRaw(SDB_STREAM, MND_STREAM_VER_NUMBER, size);
  if (pRaw == NULL) goto STREAM_ENCODE_OVER;

wafwerar's avatar
wafwerar 已提交
92
  buf = taosMemoryMalloc(tlen);
L
Liu Jicong 已提交
93 94
  if (buf == NULL) goto STREAM_ENCODE_OVER;

H
Hongze Cheng 已提交
95
  tEncoderInit(&encoder, buf, tlen);
L
Liu Jicong 已提交
96
  if (tEncodeSStreamObj(&encoder, pStream) < 0) {
H
Hongze Cheng 已提交
97
    tEncoderClear(&encoder);
L
Liu Jicong 已提交
98 99
    goto STREAM_ENCODE_OVER;
  }
H
Hongze Cheng 已提交
100
  tEncoderClear(&encoder);
L
Liu Jicong 已提交
101 102 103 104 105 106 107 108 109

  int32_t dataPos = 0;
  SDB_SET_INT32(pRaw, dataPos, tlen, STREAM_ENCODE_OVER);
  SDB_SET_BINARY(pRaw, dataPos, buf, tlen, STREAM_ENCODE_OVER);
  SDB_SET_DATALEN(pRaw, dataPos, STREAM_ENCODE_OVER);

  terrno = TSDB_CODE_SUCCESS;

STREAM_ENCODE_OVER:
wafwerar's avatar
wafwerar 已提交
110
  taosMemoryFreeClear(buf);
L
Liu Jicong 已提交
111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142
  if (terrno != TSDB_CODE_SUCCESS) {
    mError("stream:%s, failed to encode to raw:%p since %s", pStream->name, pRaw, terrstr());
    sdbFreeRaw(pRaw);
    return NULL;
  }

  mTrace("stream:%s, encode to raw:%p, row:%p", pStream->name, pRaw, pStream);
  return pRaw;
}

SSdbRow *mndStreamActionDecode(SSdbRaw *pRaw) {
  terrno = TSDB_CODE_OUT_OF_MEMORY;
  void *buf = NULL;

  int8_t sver = 0;
  if (sdbGetRawSoftVer(pRaw, &sver) != 0) goto STREAM_DECODE_OVER;

  if (sver != MND_STREAM_VER_NUMBER) {
    terrno = TSDB_CODE_SDB_INVALID_DATA_VER;
    goto STREAM_DECODE_OVER;
  }

  int32_t  size = sizeof(SStreamObj);
  SSdbRow *pRow = sdbAllocRow(size);
  if (pRow == NULL) goto STREAM_DECODE_OVER;

  SStreamObj *pStream = sdbGetRowObj(pRow);
  if (pStream == NULL) goto STREAM_DECODE_OVER;

  int32_t tlen;
  int32_t dataPos = 0;
  SDB_GET_INT32(pRaw, dataPos, &tlen, STREAM_DECODE_OVER);
wafwerar's avatar
wafwerar 已提交
143
  buf = taosMemoryMalloc(tlen + 1);
L
Liu Jicong 已提交
144 145 146
  if (buf == NULL) goto STREAM_DECODE_OVER;
  SDB_GET_BINARY(pRaw, dataPos, buf, tlen, STREAM_DECODE_OVER);

H
Hongze Cheng 已提交
147 148
  SDecoder decoder;
  tDecoderInit(&decoder, buf, tlen + 1);
L
Liu Jicong 已提交
149
  if (tDecodeSStreamObj(&decoder, pStream) < 0) {
L
Liu Jicong 已提交
150
    tDecoderClear(&decoder);
L
Liu Jicong 已提交
151 152
    goto STREAM_DECODE_OVER;
  }
L
Liu Jicong 已提交
153
  tDecoderClear(&decoder);
L
Liu Jicong 已提交
154 155 156 157

  terrno = TSDB_CODE_SUCCESS;

STREAM_DECODE_OVER:
wafwerar's avatar
wafwerar 已提交
158
  taosMemoryFreeClear(buf);
L
Liu Jicong 已提交
159 160
  if (terrno != TSDB_CODE_SUCCESS) {
    mError("stream:%s, failed to decode from raw:%p since %s", pStream->name, pRaw, terrstr());
wafwerar's avatar
wafwerar 已提交
161
    taosMemoryFreeClear(pRow);
L
Liu Jicong 已提交
162 163 164 165 166 167 168 169 170 171 172 173 174 175
    return NULL;
  }

  mTrace("stream:%s, decode from raw:%p, row:%p", pStream->name, pRaw, pStream);
  return pRow;
}

static int32_t mndStreamActionInsert(SSdb *pSdb, SStreamObj *pStream) {
  mTrace("stream:%s, perform insert action", pStream->name);
  return 0;
}

static int32_t mndStreamActionDelete(SSdb *pSdb, SStreamObj *pStream) {
  mTrace("stream:%s, perform delete action", pStream->name);
L
Liu Jicong 已提交
176 177 178
  taosWLockLatch(&pStream->lock);
  tFreeStreamObj(pStream);
  taosWUnLockLatch(&pStream->lock);
L
Liu Jicong 已提交
179 180 181 182 183
  return 0;
}

static int32_t mndStreamActionUpdate(SSdb *pSdb, SStreamObj *pOldStream, SStreamObj *pNewStream) {
  mTrace("stream:%s, perform update action", pOldStream->name);
wafwerar's avatar
wafwerar 已提交
184
  atomic_exchange_64(&pOldStream->updateTime, pNewStream->updateTime);
L
Liu Jicong 已提交
185 186 187 188
  atomic_exchange_32(&pOldStream->version, pNewStream->version);

  taosWLockLatch(&pOldStream->lock);

189
  pOldStream->status = pNewStream->status;
L
Liu Jicong 已提交
190 191 192 193 194 195 196 197 198 199 200 201 202 203 204 205 206 207 208

  taosWUnLockLatch(&pOldStream->lock);
  return 0;
}

SStreamObj *mndAcquireStream(SMnode *pMnode, char *streamName) {
  SSdb       *pSdb = pMnode->pSdb;
  SStreamObj *pStream = sdbAcquire(pSdb, SDB_STREAM, streamName);
  if (pStream == NULL && terrno == TSDB_CODE_SDB_OBJ_NOT_THERE) {
    terrno = TSDB_CODE_MND_STREAM_NOT_EXIST;
  }
  return pStream;
}

void mndReleaseStream(SMnode *pMnode, SStreamObj *pStream) {
  SSdb *pSdb = pMnode->pSdb;
  sdbRelease(pSdb, pStream);
}

L
Liu Jicong 已提交
209 210 211 212 213 214 215 216 217 218 219 220 221 222 223 224 225 226 227 228 229 230 231 232
static void mndShowStreamStatus(char *dst, SStreamObj *pStream) {
  int8_t status = atomic_load_8(&pStream->status);
  if (status == STREAM_STATUS__NORMAL) {
    strcpy(dst, "normal");
  } else if (status == STREAM_STATUS__STOP) {
    strcpy(dst, "stop");
  } else if (status == STREAM_STATUS__FAILED) {
    strcpy(dst, "failed");
  } else if (status == STREAM_STATUS__RECOVER) {
    strcpy(dst, "recover");
  }
}

static void mndShowStreamTrigger(char *dst, SStreamObj *pStream) {
  int8_t trigger = pStream->trigger;
  if (trigger == STREAM_TRIGGER_AT_ONCE) {
    strcpy(dst, "at once");
  } else if (trigger == STREAM_TRIGGER_WINDOW_CLOSE) {
    strcpy(dst, "window close");
  } else if (trigger == STREAM_TRIGGER_MAX_DELAY) {
    strcpy(dst, "max delay");
  }
}

L
Liu Jicong 已提交
233
static int32_t mndCheckCreateStreamReq(SCMCreateStreamReq *pCreate) {
L
Liu Jicong 已提交
234 235
  if (pCreate->name[0] == 0 || pCreate->sql == NULL || pCreate->sql[0] == 0 || pCreate->sourceDB[0] == 0 ||
      pCreate->targetStbFullName[0] == 0) {
L
Liu Jicong 已提交
236 237 238 239 240 241
    terrno = TSDB_CODE_MND_INVALID_STREAM_OPTION;
    return -1;
  }
  return 0;
}

X
Xiaoyu Wang 已提交
242
static int32_t mndStreamGetPlanString(const char *ast, int8_t triggerType, int64_t watermark, char **pStr) {
S
sma  
Shengliang Guan 已提交
243
  if (NULL == ast) {
L
Liu Jicong 已提交
244 245 246 247
    return TSDB_CODE_SUCCESS;
  }

  SNode  *pAst = NULL;
S
sma  
Shengliang Guan 已提交
248
  int32_t code = nodesStringToNode(ast, &pAst);
L
Liu Jicong 已提交
249 250 251 252 253 254 255

  SQueryPlan *pPlan = NULL;
  if (TSDB_CODE_SUCCESS == code) {
    SPlanContext cxt = {
        .pAstRoot = pAst,
        .topicQuery = false,
        .streamQuery = true,
256
        .triggerType = triggerType == STREAM_TRIGGER_MAX_DELAY ? STREAM_TRIGGER_WINDOW_CLOSE : triggerType,
X
Xiaoyu Wang 已提交
257
        .watermark = watermark,
L
Liu Jicong 已提交
258 259 260 261 262
    };
    code = qCreateQueryPlan(&cxt, &pPlan, NULL);
  }

  if (TSDB_CODE_SUCCESS == code) {
263
    code = nodesNodeToString((SNode *)pPlan, false, pStr, NULL);
L
Liu Jicong 已提交
264 265
  }
  nodesDestroyNode(pAst);
266
  nodesDestroyNode((SNode *)pPlan);
L
Liu Jicong 已提交
267 268 269 270
  terrno = code;
  return code;
}

L
Liu Jicong 已提交
271
static int32_t mndBuildStreamObjFromCreateReq(SMnode *pMnode, SStreamObj *pObj, SCMCreateStreamReq *pCreate) {
272 273
  SNode      *pAst = NULL;
  SQueryPlan *pPlan = NULL;
L
Liu Jicong 已提交
274

275
  mInfo("stream:%s to create", pCreate->name);
L
Liu Jicong 已提交
276 277 278 279 280 281 282 283 284
  memcpy(pObj->name, pCreate->name, TSDB_STREAM_FNAME_LEN);
  pObj->createTime = taosGetTimestampMs();
  pObj->updateTime = pObj->createTime;
  pObj->version = 1;
  pObj->smaId = 0;

  pObj->uid = mndGenerateUid(pObj->name, strlen(pObj->name));
  pObj->status = 0;

285
  pObj->igExpired = pCreate->igExpired;
L
Liu Jicong 已提交
286 287 288
  pObj->trigger = pCreate->triggerType;
  pObj->triggerParam = pCreate->maxDelay;
  pObj->watermark = pCreate->watermark;
L
Liu Jicong 已提交
289
  pObj->fillHistory = pCreate->fillHistory;
L
Liu Jicong 已提交
290 291 292 293

  memcpy(pObj->sourceDb, pCreate->sourceDB, TSDB_DB_FNAME_LEN);
  SDbObj *pSourceDb = mndAcquireDb(pMnode, pCreate->sourceDB);
  if (pSourceDb == NULL) {
294
    mInfo("stream:%s failed to create, source db %s not exist since %s", pCreate->name, pObj->sourceDb, terrstr());
L
Liu Jicong 已提交
295 296 297
    return -1;
  }
  pObj->sourceDbUid = pSourceDb->uid;
298
  mndReleaseDb(pMnode, pSourceDb);
L
Liu Jicong 已提交
299 300 301 302 303

  memcpy(pObj->targetSTbName, pCreate->targetStbFullName, TSDB_TABLE_FNAME_LEN);

  SDbObj *pTargetDb = mndAcquireDbByStb(pMnode, pObj->targetSTbName);
  if (pTargetDb == NULL) {
304
    mInfo("stream:%s failed to create, target db %s not exist since %s", pCreate->name, pObj->targetDb, terrstr());
L
Liu Jicong 已提交
305 306 307 308 309 310
    return -1;
  }
  tstrncpy(pObj->targetDb, pTargetDb->name, TSDB_DB_FNAME_LEN);

  pObj->targetStbUid = mndGenerateUid(pObj->targetSTbName, TSDB_TABLE_FNAME_LEN);
  pObj->targetDbUid = pTargetDb->uid;
311
  mndReleaseDb(pMnode, pTargetDb);
L
Liu Jicong 已提交
312 313 314 315 316 317 318 319 320 321 322 323 324 325 326 327 328 329 330 331 332 333 334 335 336

  pObj->sql = pCreate->sql;
  pObj->ast = pCreate->ast;

  pCreate->sql = NULL;
  pCreate->ast = NULL;

  // deserialize ast
  if (nodesStringToNode(pObj->ast, &pAst) < 0) {
    /*ASSERT(0);*/
    goto FAIL;
  }

  // extract output schema from ast
  if (qExtractResultSchema(pAst, (int32_t *)&pObj->outputSchema.nCols, &pObj->outputSchema.pSchema) != 0) {
    /*ASSERT(0);*/
    goto FAIL;
  }

  SPlanContext cxt = {
      .pAstRoot = pAst,
      .topicQuery = false,
      .streamQuery = true,
      .triggerType = pObj->trigger == STREAM_TRIGGER_MAX_DELAY ? STREAM_TRIGGER_WINDOW_CLOSE : pObj->trigger,
      .watermark = pObj->watermark,
337
      .igExpired = pObj->igExpired,
L
Liu Jicong 已提交
338 339 340 341 342 343 344 345 346 347 348 349 350 351
  };

  // using ast and param to build physical plan
  if (qCreateQueryPlan(&cxt, &pPlan, NULL) < 0) {
    /*ASSERT(0);*/
    goto FAIL;
  }

  // save physcial plan
  if (nodesNodeToString((SNode *)pPlan, false, &pObj->physicalPlan, NULL) != 0) {
    /*ASSERT(0);*/
    goto FAIL;
  }

352 353 354 355 356 357 358 359 360 361 362 363 364 365
  pObj->tagSchema.nCols = pCreate->numOfTags;
  if (pCreate->numOfTags) {
    pObj->tagSchema.pSchema = taosMemoryCalloc(pCreate->numOfTags, sizeof(SSchema));
  }
  ASSERT(pCreate->numOfTags == taosArrayGetSize(pCreate->pTags));
  for (int32_t i = 0; i < pCreate->numOfTags; i++) {
    SField *pField = taosArrayGet(pCreate->pTags, i);
    pObj->tagSchema.pSchema[i].colId = pObj->outputSchema.nCols + i + 1;
    pObj->tagSchema.pSchema[i].bytes = pField->bytes;
    pObj->tagSchema.pSchema[i].flags = pField->flags;
    pObj->tagSchema.pSchema[i].type = pField->type;
    memcpy(pObj->tagSchema.pSchema[i].name, pField->name, TSDB_COL_NAME_LEN);
  }

L
Liu Jicong 已提交
366 367
FAIL:
  if (pAst != NULL) nodesDestroyNode(pAst);
368
  if (pPlan != NULL) qDestroyQueryPlan(pPlan);
L
Liu Jicong 已提交
369 370 371
  return 0;
}

372
int32_t mndPersistTaskDeployReq(STrans *pTrans, const SStreamTask *pTask) {
373
  if (pTask->taskLevel == TASK_LEVEL__AGG) {
L
Liu Jicong 已提交
374 375
    ASSERT(taosArrayGetSize(pTask->childEpInfo) != 0);
  }
376 377 378 379 380 381 382 383 384
  SEncoder encoder;
  tEncoderInit(&encoder, NULL, 0);
  tEncodeSStreamTask(&encoder, pTask);
  int32_t size = encoder.pos;
  int32_t tlen = sizeof(SMsgHead) + size;
  tEncoderClear(&encoder);
  void *buf = taosMemoryCalloc(1, tlen);
  if (buf == NULL) {
    terrno = TSDB_CODE_OUT_OF_MEMORY;
L
Liu Jicong 已提交
385 386
    return -1;
  }
387 388 389 390 391
  ((SMsgHead *)buf)->vgId = htonl(pTask->nodeId);
  void *abuf = POINTER_SHIFT(buf, sizeof(SMsgHead));
  tEncoderInit(&encoder, abuf, size);
  tEncodeSStreamTask(&encoder, pTask);
  tEncoderClear(&encoder);
L
Liu Jicong 已提交
392

393 394 395 396 397 398 399
  STransAction action = {0};
  memcpy(&action.epSet, &pTask->epSet, sizeof(SEpSet));
  action.pCont = buf;
  action.contLen = tlen;
  action.msgType = TDMT_STREAM_TASK_DEPLOY;
  if (mndTransAppendRedoAction(pTrans, &action) != 0) {
    taosMemoryFree(buf);
L
Liu Jicong 已提交
400 401 402 403 404
    return -1;
  }
  return 0;
}

405 406 407 408 409 410 411 412 413 414 415
int32_t mndPersistStreamTasks(SMnode *pMnode, STrans *pTrans, SStreamObj *pStream) {
  int32_t level = taosArrayGetSize(pStream->tasks);
  for (int32_t i = 0; i < level; i++) {
    SArray *pLevel = taosArrayGetP(pStream->tasks, i);
    int32_t sz = taosArrayGetSize(pLevel);
    for (int32_t j = 0; j < sz; j++) {
      SStreamTask *pTask = taosArrayGetP(pLevel, j);
      if (mndPersistTaskDeployReq(pTrans, pTask) < 0) {
        return -1;
      }
    }
L
Liu Jicong 已提交
416
  }
417 418
  return 0;
}
L
Liu Jicong 已提交
419

420 421
int32_t mndPersistStream(SMnode *pMnode, STrans *pTrans, SStreamObj *pStream) {
  if (mndPersistStreamTasks(pMnode, pTrans, pStream) < 0) {
L
Liu Jicong 已提交
422 423
    return -1;
  }
424 425 426
  SSdbRaw *pCommitRaw = mndStreamActionEncode(pStream);
  if (pCommitRaw == NULL || mndTransAppendCommitlog(pTrans, pCommitRaw) != 0) {
    mError("trans:%d, failed to append commit log since %s", pTrans->id, terrstr());
L
Liu Jicong 已提交
427 428
    return -1;
  }
S
Shengliang Guan 已提交
429
  (void)sdbSetRawStatus(pCommitRaw, SDB_STATUS_READY);
430 431
  return 0;
}
L
Liu Jicong 已提交
432

433
int32_t mndPersistDropStreamLog(SMnode *pMnode, STrans *pTrans, SStreamObj *pStream) {
434 435 436
  SSdbRaw *pCommitRaw = mndStreamActionEncode(pStream);
  if (pCommitRaw == NULL || mndTransAppendCommitlog(pTrans, pCommitRaw) != 0) {
    mError("trans:%d, failed to append commit log since %s", pTrans->id, terrstr());
L
Liu Jicong 已提交
437 438 439
    mndTransDrop(pTrans);
    return -1;
  }
S
Shengliang Guan 已提交
440
  (void)sdbSetRawStatus(pCommitRaw, SDB_STATUS_DROPPED);
S
sma  
Shengliang Guan 已提交
441 442 443
  return 0;
}

444 445 446 447
static int32_t mndSetStreamRecover(SMnode *pMnode, STrans *pTrans, const SStreamObj *pStream) {
  SStreamObj streamObj = {0};
  memcpy(streamObj.name, pStream->name, TSDB_STREAM_FNAME_LEN);
  streamObj.status = STREAM_STATUS__RECOVER;
S
Shengliang Guan 已提交
448

449
  SSdbRaw *pCommitRaw = mndStreamActionEncode(&streamObj);
S
Shengliang Guan 已提交
450 451
  if (pCommitRaw == NULL) return -1;
  if (mndTransAppendCommitlog(pTrans, pCommitRaw) != 0) {
452 453 454 455
    mError("stream trans:%d, failed to append commit log since %s", pTrans->id, terrstr());
    mndTransDrop(pTrans);
    return -1;
  }
S
Shengliang Guan 已提交
456
  (void)sdbSetRawStatus(pCommitRaw, SDB_STATUS_READY);
457 458 459
  return 0;
}

L
Liu Jicong 已提交
460
static int32_t mndCreateStbForStream(SMnode *pMnode, STrans *pTrans, const SStreamObj *pStream, const char *user) {
L
Liu Jicong 已提交
461 462
  SStbObj *pStb = NULL;
  SDbObj  *pDb = NULL;
463 464 465 466 467 468 469

  SMCreateStbReq createReq = {0};
  tstrncpy(createReq.name, pStream->targetSTbName, TSDB_TABLE_FNAME_LEN);
  createReq.numOfColumns = pStream->outputSchema.nCols;
  createReq.numOfTags = 1;  // group id
  createReq.pColumns = taosArrayInit(createReq.numOfColumns, sizeof(SField));
  // build fields
L
Liu Jicong 已提交
470 471 472 473 474 475 476 477 478 479
  taosArraySetSize(createReq.pColumns, createReq.numOfColumns);
  for (int32_t i = 0; i < createReq.numOfColumns; i++) {
    SField *pField = taosArrayGet(createReq.pColumns, i);
    tstrncpy(pField->name, pStream->outputSchema.pSchema[i].name, TSDB_COL_NAME_LEN);
    pField->flags = pStream->outputSchema.pSchema[i].flags;
    pField->type = pStream->outputSchema.pSchema[i].type;
    pField->bytes = pStream->outputSchema.pSchema[i].bytes;
  }
  createReq.pTags = taosArrayInit(createReq.numOfTags, sizeof(SField));
  taosArraySetSize(createReq.pTags, 1);
480
  // build tags
L
Liu Jicong 已提交
481 482 483 484 485
  SField *pField = taosArrayGet(createReq.pTags, 0);
  strcpy(pField->name, "group_id");
  pField->type = TSDB_DATA_TYPE_UBIGINT;
  pField->flags = 0;
  pField->bytes = 8;
486 487 488 489 490 491 492 493 494 495 496 497 498 499 500 501 502 503

  if (mndCheckCreateStbReq(&createReq) != 0) {
    goto _OVER;
  }

  pStb = mndAcquireStb(pMnode, createReq.name);
  if (pStb != NULL) {
    terrno = TSDB_CODE_MND_STB_ALREADY_EXIST;
    goto _OVER;
  }

  pDb = mndAcquireDbByStb(pMnode, createReq.name);
  if (pDb == NULL) {
    terrno = TSDB_CODE_MND_DB_NOT_SELECTED;
    goto _OVER;
  }

  int32_t numOfStbs = -1;
S
Shengliang Guan 已提交
504 505 506 507
  if (mndGetNumOfStbs(pMnode, pDb->name, &numOfStbs) != 0) {
    goto _OVER;
  }

508 509 510 511 512 513 514 515 516 517 518
  if (pDb->cfg.numOfStables == 1 && numOfStbs != 0) {
    terrno = TSDB_CODE_MND_SINGLE_STB_MODE_DB;
    goto _OVER;
  }

  SStbObj stbObj = {0};

  if (mndBuildStbFromReq(pMnode, &stbObj, &createReq, pDb) != 0) {
    goto _OVER;
  }

L
Liu Jicong 已提交
519 520
  stbObj.uid = pStream->targetStbUid;

L
Liu Jicong 已提交
521 522 523 524 525 526 527
  if (mndAddStbToTrans(pMnode, pTrans, pDb, &stbObj) < 0) {
    mndFreeStb(&stbObj);
    goto _OVER;
  }

  tFreeSMCreateStbReq(&createReq);
  mndFreeStb(&stbObj);
528
  mndReleaseDb(pMnode, pDb);
529

L
Liu Jicong 已提交
530
  return 0;
531
_OVER:
L
Liu Jicong 已提交
532
  tFreeSMCreateStbReq(&createReq);
533 534
  mndReleaseStb(pMnode, pStb);
  mndReleaseDb(pMnode, pDb);
L
Liu Jicong 已提交
535
  return -1;
536 537
}

L
Liu Jicong 已提交
538 539 540 541
static int32_t mndPersistTaskDropReq(STrans *pTrans, SStreamTask *pTask) {
  ASSERT(pTask->nodeId != 0);

  // vnode
L
Liu Jicong 已提交
542 543 544 545 546 547 548 549 550 551 552 553
  /*if (pTask->nodeId > 0) {*/
  SVDropStreamTaskReq *pReq = taosMemoryCalloc(1, sizeof(SVDropStreamTaskReq));
  if (pReq == NULL) {
    terrno = TSDB_CODE_OUT_OF_MEMORY;
    return -1;
  }
  pReq->head.vgId = htonl(pTask->nodeId);
  pReq->taskId = pTask->taskId;
  STransAction action = {0};
  memcpy(&action.epSet, &pTask->epSet, sizeof(SEpSet));
  action.pCont = pReq;
  action.contLen = sizeof(SVDropStreamTaskReq);
L
Liu Jicong 已提交
554
  action.msgType = TDMT_STREAM_TASK_DROP;
L
Liu Jicong 已提交
555 556 557
  if (mndTransAppendRedoAction(pTrans, &action) != 0) {
    taosMemoryFree(pReq);
    return -1;
L
Liu Jicong 已提交
558
  }
L
Liu Jicong 已提交
559
  /*}*/
L
Liu Jicong 已提交
560 561 562 563

  return 0;
}

L
Liu Jicong 已提交
564
int32_t mndDropStreamTasks(SMnode *pMnode, STrans *pTrans, SStreamObj *pStream) {
L
Liu Jicong 已提交
565 566 567 568 569 570 571 572 573 574 575 576 577 578
  int32_t lv = taosArrayGetSize(pStream->tasks);
  for (int32_t i = 0; i < lv; i++) {
    SArray *pTasks = taosArrayGetP(pStream->tasks, i);
    int32_t sz = taosArrayGetSize(pTasks);
    for (int32_t j = 0; j < sz; j++) {
      SStreamTask *pTask = taosArrayGetP(pTasks, j);
      if (mndPersistTaskDropReq(pTrans, pTask) < 0) {
        return -1;
      }
    }
  }
  return 0;
}

S
Shengliang Guan 已提交
579 580
static int32_t mndProcessCreateStreamReq(SRpcMsg *pReq) {
  SMnode            *pMnode = pReq->info.node;
L
Liu Jicong 已提交
581 582 583 584
  int32_t            code = -1;
  SStreamObj        *pStream = NULL;
  SDbObj            *pDb = NULL;
  SCMCreateStreamReq createStreamReq = {0};
585
  SStreamObj         streamObj = {0};
L
Liu Jicong 已提交
586

S
Shengliang Guan 已提交
587
  if (tDeserializeSCMCreateStreamReq(pReq->pCont, pReq->contLen, &createStreamReq) != 0) {
L
Liu Jicong 已提交
588
    terrno = TSDB_CODE_INVALID_MSG;
589
    goto _OVER;
L
Liu Jicong 已提交
590 591
  }

592
  mInfo("stream:%s, start to create, sql:%s", createStreamReq.name, createStreamReq.sql);
L
Liu Jicong 已提交
593 594 595

  if (mndCheckCreateStreamReq(&createStreamReq) != 0) {
    mError("stream:%s, failed to create since %s", createStreamReq.name, terrstr());
596
    goto _OVER;
L
Liu Jicong 已提交
597 598 599 600 601
  }

  pStream = mndAcquireStream(pMnode, createStreamReq.name);
  if (pStream != NULL) {
    if (createStreamReq.igExists) {
602
      mInfo("stream:%s, already exist, ignore exist is set", createStreamReq.name);
L
Liu Jicong 已提交
603
      code = 0;
604
      goto _OVER;
L
Liu Jicong 已提交
605 606
    } else {
      terrno = TSDB_CODE_MND_STREAM_ALREADY_EXIST;
607
      goto _OVER;
L
Liu Jicong 已提交
608 609
    }
  } else if (terrno != TSDB_CODE_MND_STREAM_NOT_EXIST) {
610
    goto _OVER;
L
Liu Jicong 已提交
611 612
  }

L
Liu Jicong 已提交
613 614 615 616 617 618
  // build stream obj from request
  if (mndBuildStreamObjFromCreateReq(pMnode, &streamObj, &createStreamReq) < 0) {
    mError("stream:%s, failed to create since %s", createStreamReq.name, terrstr());
    goto _OVER;
  }

619
  STrans *pTrans = mndTransCreate(pMnode, TRN_POLICY_ROLLBACK, TRN_CONFLICT_DB_INSIDE, pReq, "create-stream");
620 621 622 623
  if (pTrans == NULL) {
    mError("stream:%s, failed to create since %s", createStreamReq.name, terrstr());
    goto _OVER;
  }
L
Liu Jicong 已提交
624
  mndTransSetDbName(pTrans, createStreamReq.sourceDB, streamObj.targetDb);
625
  mInfo("trans:%d, used to create stream:%s", pTrans->id, createStreamReq.name);
626

L
Liu Jicong 已提交
627
  // create stb for stream
628
  if (mndCreateStbForStream(pMnode, pTrans, &streamObj, pReq->info.conn.user) < 0) {
L
Liu Jicong 已提交
629 630 631 632 633
    mError("trans:%d, failed to create stb for stream %s since %s", pTrans->id, createStreamReq.name, terrstr());
    mndTransDrop(pTrans);
    goto _OVER;
  }

L
Liu Jicong 已提交
634
  // schedule stream task for stream obj
635
  if (mndScheduleStream(pMnode, &streamObj) < 0) {
L
Liu Jicong 已提交
636 637 638 639 640
    mError("stream:%s, failed to schedule since %s", createStreamReq.name, terrstr());
    mndTransDrop(pTrans);
    goto _OVER;
  }

L
Liu Jicong 已提交
641
  // add stream to trans
L
Liu Jicong 已提交
642 643 644 645 646 647
  if (mndPersistStream(pMnode, pTrans, &streamObj) < 0) {
    mError("stream:%s, failed to schedule since %s", createStreamReq.name, terrstr());
    mndTransDrop(pTrans);
    goto _OVER;
  }

648 649 650 651 652 653 654 655 656 657
  if (mndCheckDbPrivilegeByName(pMnode, pReq->info.conn.user, MND_OPER_READ_DB, streamObj.sourceDb) != 0) {
    mndTransDrop(pTrans);
    goto _OVER;
  }

  if (mndCheckDbPrivilegeByName(pMnode, pReq->info.conn.user, MND_OPER_WRITE_DB, streamObj.targetDb) != 0) {
    mndTransDrop(pTrans);
    goto _OVER;
  }

L
Liu Jicong 已提交
658
  // execute creation
L
Liu Jicong 已提交
659 660 661 662 663
  if (mndTransPrepare(pMnode, pTrans) != 0) {
    mError("trans:%d, failed to prepare since %s", pTrans->id, terrstr());
    mndTransDrop(pTrans);
    goto _OVER;
  }
L
Liu Jicong 已提交
664

L
Liu Jicong 已提交
665 666 667
  mndTransDrop(pTrans);

  code = TSDB_CODE_ACTION_IN_PROGRESS;
L
Liu Jicong 已提交
668

669
_OVER:
S
Shengliang Guan 已提交
670
  if (code != 0 && code != TSDB_CODE_ACTION_IN_PROGRESS) {
L
Liu Jicong 已提交
671 672 673 674 675 676 677
    mError("stream:%s, failed to create since %s", createStreamReq.name, terrstr());
  }

  mndReleaseStream(pMnode, pStream);
  mndReleaseDb(pMnode, pDb);

  tFreeSCMCreateStreamReq(&createStreamReq);
L
Liu Jicong 已提交
678
  tFreeStreamObj(&streamObj);
L
Liu Jicong 已提交
679 680 681
  return code;
}

L
Liu Jicong 已提交
682 683 684 685 686 687
static int32_t mndProcessDropStreamReq(SRpcMsg *pReq) {
  SMnode     *pMnode = pReq->info.node;
  SStreamObj *pStream = NULL;
  /*SDbObj     *pDb = NULL;*/
  /*SUserObj   *pUser = NULL;*/

L
Liu Jicong 已提交
688 689 690 691 692 693
  SMDropStreamReq dropReq = {0};
  if (tDeserializeSMDropStreamReq(pReq->pCont, pReq->contLen, &dropReq) < 0) {
    ASSERT(0);
    terrno = TSDB_CODE_INVALID_MSG;
    return -1;
  }
L
Liu Jicong 已提交
694

L
Liu Jicong 已提交
695
  pStream = mndAcquireStream(pMnode, dropReq.name);
L
Liu Jicong 已提交
696 697

  if (pStream == NULL) {
L
Liu Jicong 已提交
698
    if (dropReq.igNotExists) {
699
      mInfo("stream:%s, not exist, ignore not exist is set", dropReq.name);
L
Liu Jicong 已提交
700
      sdbRelease(pMnode->pSdb, pStream);
701
      return 0;
L
Liu Jicong 已提交
702 703 704 705
    } else {
      terrno = TSDB_CODE_MND_STREAM_NOT_EXIST;
      return -1;
    }
L
Liu Jicong 已提交
706 707
  }

708 709
  if (mndCheckDbPrivilegeByName(pMnode, pReq->info.conn.user, MND_OPER_WRITE_DB, pStream->targetDb) != 0) {
    return -1;
L
Liu Jicong 已提交
710 711
  }

712 713
  STrans *pTrans = mndTransCreate(pMnode, TRN_POLICY_RETRY, TRN_CONFLICT_DB_INSIDE, pReq, "drop-stream");
  mndTransSetDbName(pTrans, pStream->sourceDb, pStream->targetDb);
L
Liu Jicong 已提交
714
  if (pTrans == NULL) {
L
Liu Jicong 已提交
715
    mError("stream:%s, failed to drop since %s", dropReq.name, terrstr());
L
Liu Jicong 已提交
716 717
    sdbRelease(pMnode->pSdb, pStream);
    return -1;
L
Liu Jicong 已提交
718
  }
719
  mInfo("trans:%d, used to drop stream:%s", pTrans->id, dropReq.name);
L
Liu Jicong 已提交
720 721 722

  // drop all tasks
  if (mndDropStreamTasks(pMnode, pTrans, pStream) < 0) {
L
Liu Jicong 已提交
723
    mError("stream:%s, failed to drop task since %s", dropReq.name, terrstr());
L
Liu Jicong 已提交
724
    sdbRelease(pMnode->pSdb, pStream);
S
Shengliang Guan 已提交
725
    mndTransDrop(pTrans);
L
Liu Jicong 已提交
726
    return -1;
L
Liu Jicong 已提交
727 728
  }

L
Liu Jicong 已提交
729
  // drop stream
L
Liu Jicong 已提交
730 731
  if (mndPersistDropStreamLog(pMnode, pTrans, pStream) < 0) {
    sdbRelease(pMnode->pSdb, pStream);
S
Shengliang Guan 已提交
732
    mndTransDrop(pTrans);
L
Liu Jicong 已提交
733 734 735
    return -1;
  }

L
Liu Jicong 已提交
736 737 738 739 740 741 742 743
  if (mndTransPrepare(pMnode, pTrans) != 0) {
    mError("trans:%d, failed to prepare drop stream trans since %s", pTrans->id, terrstr());
    sdbRelease(pMnode->pSdb, pStream);
    mndTransDrop(pTrans);
    return -1;
  }

  sdbRelease(pMnode->pSdb, pStream);
L
Liu Jicong 已提交
744
  mndTransDrop(pTrans);
L
Liu Jicong 已提交
745 746

  return TSDB_CODE_ACTION_IN_PROGRESS;
L
Liu Jicong 已提交
747 748
}

L
Liu Jicong 已提交
749
#if 0
L
Liu Jicong 已提交
750 751 752 753 754 755 756 757 758 759 760 761 762 763 764 765 766
static int32_t mndProcessRecoverStreamReq(SRpcMsg *pReq) {
  SMnode     *pMnode = pReq->info.node;
  SStreamObj *pStream = NULL;
  /*SDbObj     *pDb = NULL;*/
  /*SUserObj   *pUser = NULL;*/

  SMRecoverStreamReq recoverReq = {0};
  if (tDeserializeSMRecoverStreamReq(pReq->pCont, pReq->contLen, &recoverReq) < 0) {
    ASSERT(0);
    terrno = TSDB_CODE_INVALID_MSG;
    return -1;
  }

  pStream = mndAcquireStream(pMnode, recoverReq.name);

  if (pStream == NULL) {
    if (recoverReq.igNotExists) {
767
      mInfo("stream:%s, not exist, ignore not exist is set", recoverReq.name);
L
Liu Jicong 已提交
768 769 770 771 772 773 774 775 776 777 778 779 780 781 782 783 784 785
      sdbRelease(pMnode->pSdb, pStream);
      return 0;
    } else {
      terrno = TSDB_CODE_MND_STREAM_NOT_EXIST;
      return -1;
    }
  }

  if (mndCheckDbPrivilegeByName(pMnode, pReq->info.conn.user, MND_OPER_WRITE_DB, pStream->targetDb) != 0) {
    return -1;
  }

  STrans *pTrans = mndTransCreate(pMnode, TRN_POLICY_RETRY, TRN_CONFLICT_NOTHING, pReq);
  if (pTrans == NULL) {
    mError("stream:%s, failed to recover since %s", recoverReq.name, terrstr());
    sdbRelease(pMnode->pSdb, pStream);
    return -1;
  }
786
  mInfo("trans:%d, used to drop stream:%s", pTrans->id, recoverReq.name);
L
Liu Jicong 已提交
787 788

  // broadcast to recover all tasks
789
  if (mndRecoverStreamTasks(pMnode, pTrans, pStream) < 0) {
L
Liu Jicong 已提交
790 791 792 793 794 795
    mError("stream:%s, failed to recover task since %s", recoverReq.name, terrstr());
    sdbRelease(pMnode->pSdb, pStream);
    return -1;
  }

  // update stream status
796
  if (mndSetStreamRecover(pMnode, pTrans, pStream) < 0) {
L
Liu Jicong 已提交
797 798 799 800 801 802 803 804 805 806 807 808 809 810 811
    sdbRelease(pMnode->pSdb, pStream);
    return -1;
  }

  if (mndTransPrepare(pMnode, pTrans) != 0) {
    mError("trans:%d, failed to prepare recover stream trans since %s", pTrans->id, terrstr());
    sdbRelease(pMnode->pSdb, pStream);
    mndTransDrop(pTrans);
    return -1;
  }

  sdbRelease(pMnode->pSdb, pStream);

  return TSDB_CODE_ACTION_IN_PROGRESS;
}
L
Liu Jicong 已提交
812
#endif
L
Liu Jicong 已提交
813

L
Liu Jicong 已提交
814 815
int32_t mndDropStreamByDb(SMnode *pMnode, STrans *pTrans, SDbObj *pDb) {
  SSdb *pSdb = pMnode->pSdb;
S
Shengliang Guan 已提交
816
  void *pIter = NULL;
L
Liu Jicong 已提交
817 818

  while (1) {
S
Shengliang Guan 已提交
819
    SStreamObj *pStream = NULL;
L
Liu Jicong 已提交
820 821 822 823 824 825
    pIter = sdbFetch(pSdb, SDB_STREAM, pIter, (void **)&pStream);
    if (pIter == NULL) break;

    if (pStream->sourceDbUid == pDb->uid || pStream->targetDbUid == pDb->uid) {
      if (pStream->sourceDbUid != pStream->targetDbUid) {
        sdbRelease(pSdb, pStream);
S
Shengliang Guan 已提交
826 827 828
        sdbCancelFetch(pSdb, pIter);
        mError("db:%s, failed to drop stream:%s since sourceDbUid:%" PRId64 " not match with targetDbUid:%" PRId64,
               pDb->name, pStream->name, pStream->sourceDbUid, pStream->targetDbUid);
829
        terrno = TSDB_CODE_MND_STREAM_MUST_BE_DELETED;
L
Liu Jicong 已提交
830 831
        return -1;
      } else {
832 833 834 835 836 837 838 839
#if 0
        if (mndDropStreamTasks(pMnode, pTrans, pStream) < 0) {
          mError("stream:%s, failed to drop task since %s", pStream->name, terrstr());
          sdbRelease(pMnode->pSdb, pStream);
          sdbCancelFetch(pSdb, pIter);
          return -1;
        }
#endif
L
Liu Jicong 已提交
840 841
        if (mndPersistDropStreamLog(pMnode, pTrans, pStream) < 0) {
          sdbRelease(pSdb, pStream);
S
Shengliang Guan 已提交
842
          sdbCancelFetch(pSdb, pIter);
L
Liu Jicong 已提交
843 844 845 846 847 848 849 850 851 852 853 854 855 856 857 858 859 860
          return -1;
        }
      }
    }

#if 0
    if (mndSetDropOffsetStreamLogs(pMnode, pTrans, pStream) < 0) {
      sdbRelease(pSdb, pStream);
      goto END;
    }
#endif

    sdbRelease(pSdb, pStream);
  }

  return 0;
}

L
Liu Jicong 已提交
861 862 863 864 865 866 867 868 869 870 871 872 873 874 875
static int32_t mndGetNumOfStreams(SMnode *pMnode, char *dbName, int32_t *pNumOfStreams) {
  SSdb   *pSdb = pMnode->pSdb;
  SDbObj *pDb = mndAcquireDb(pMnode, dbName);
  if (pDb == NULL) {
    terrno = TSDB_CODE_MND_DB_NOT_SELECTED;
    return -1;
  }

  int32_t numOfStreams = 0;
  void   *pIter = NULL;
  while (1) {
    SStreamObj *pStream = NULL;
    pIter = sdbFetch(pSdb, SDB_STREAM, pIter, (void **)&pStream);
    if (pIter == NULL) break;

L
Liu Jicong 已提交
876
    if (pStream->sourceDbUid == pDb->uid) {
L
Liu Jicong 已提交
877 878 879 880 881 882 883 884 885 886 887
      numOfStreams++;
    }

    sdbRelease(pSdb, pStream);
  }

  *pNumOfStreams = numOfStreams;
  mndReleaseDb(pMnode, pDb);
  return 0;
}

S
Shengliang Guan 已提交
888 889
static int32_t mndRetrieveStream(SRpcMsg *pReq, SShowObj *pShow, SSDataBlock *pBlock, int32_t rows) {
  SMnode     *pMnode = pReq->info.node;
L
Liu Jicong 已提交
890 891 892 893 894
  SSdb       *pSdb = pMnode->pSdb;
  int32_t     numOfRows = 0;
  SStreamObj *pStream = NULL;

  while (numOfRows < rows) {
L
Liu Jicong 已提交
895
    pShow->pIter = sdbFetch(pSdb, SDB_STREAM, pShow->pIter, (void **)&pStream);
L
Liu Jicong 已提交
896 897
    if (pShow->pIter == NULL) break;

L
Liu Jicong 已提交
898 899 900
    SColumnInfoData *pColInfo;
    SName            n;
    int32_t          cols = 0;
L
Liu Jicong 已提交
901

902
    char streamName[TSDB_TABLE_NAME_LEN + VARSTR_HEADER_SIZE] = {0};
S
Shengliang Guan 已提交
903
    STR_WITH_MAXSIZE_TO_VARSTR(streamName, mndGetDbStr(pStream->name), sizeof(streamName));
L
Liu Jicong 已提交
904 905
    pColInfo = taosArrayGet(pBlock->pDataBlock, cols++);
    colDataAppend(pColInfo, numOfRows, (const char *)streamName, false);
L
Liu Jicong 已提交
906

L
Liu Jicong 已提交
907 908
    pColInfo = taosArrayGet(pBlock->pDataBlock, cols++);
    colDataAppend(pColInfo, numOfRows, (const char *)&pStream->createTime, false);
L
Liu Jicong 已提交
909

L
Liu Jicong 已提交
910
    char sql[TSDB_SHOW_SQL_LEN + VARSTR_HEADER_SIZE] = {0};
S
Shengliang Guan 已提交
911
    STR_WITH_MAXSIZE_TO_VARSTR(sql, pStream->sql, sizeof(sql));
L
Liu Jicong 已提交
912 913
    pColInfo = taosArrayGet(pBlock->pDataBlock, cols++);
    colDataAppend(pColInfo, numOfRows, (const char *)sql, false);
L
Liu Jicong 已提交
914

L
Liu Jicong 已提交
915
    char status[20 + VARSTR_HEADER_SIZE] = {0};
S
Shengliang Guan 已提交
916 917 918
    char status2[20] = {0};
    mndShowStreamStatus(status2, pStream);
    STR_WITH_MAXSIZE_TO_VARSTR(status, status2, sizeof(status));
L
Liu Jicong 已提交
919
    pColInfo = taosArrayGet(pBlock->pDataBlock, cols++);
L
Liu Jicong 已提交
920
    colDataAppend(pColInfo, numOfRows, (const char *)&status, false);
L
Liu Jicong 已提交
921

922
    char sourceDB[TSDB_DB_NAME_LEN + VARSTR_HEADER_SIZE] = {0};
S
Shengliang Guan 已提交
923
    STR_WITH_MAXSIZE_TO_VARSTR(sourceDB, mndGetDbStr(pStream->sourceDb), sizeof(sourceDB));
L
Liu Jicong 已提交
924
    pColInfo = taosArrayGet(pBlock->pDataBlock, cols++);
925
    colDataAppend(pColInfo, numOfRows, (const char *)&sourceDB, false);
L
Liu Jicong 已提交
926

927
    char targetDB[TSDB_DB_NAME_LEN + VARSTR_HEADER_SIZE] = {0};
S
Shengliang Guan 已提交
928
    STR_WITH_MAXSIZE_TO_VARSTR(targetDB, mndGetDbStr(pStream->targetDb), sizeof(targetDB));
L
Liu Jicong 已提交
929
    pColInfo = taosArrayGet(pBlock->pDataBlock, cols++);
930
    colDataAppend(pColInfo, numOfRows, (const char *)&targetDB, false);
L
Liu Jicong 已提交
931

932 933 934 935 936
    if (pStream->targetSTbName[0] == 0) {
      pColInfo = taosArrayGet(pBlock->pDataBlock, cols++);
      colDataAppend(pColInfo, numOfRows, NULL, true);
    } else {
      char targetSTB[TSDB_TABLE_NAME_LEN + VARSTR_HEADER_SIZE] = {0};
S
Shengliang Guan 已提交
937
      STR_WITH_MAXSIZE_TO_VARSTR(targetSTB, mndGetStbStr(pStream->targetSTbName), sizeof(targetSTB));
938 939 940
      pColInfo = taosArrayGet(pBlock->pDataBlock, cols++);
      colDataAppend(pColInfo, numOfRows, (const char *)&targetSTB, false);
    }
L
Liu Jicong 已提交
941 942

    pColInfo = taosArrayGet(pBlock->pDataBlock, cols++);
943
    colDataAppend(pColInfo, numOfRows, (const char *)&pStream->watermark, false);
L
Liu Jicong 已提交
944

L
Liu Jicong 已提交
945
    char trigger[20 + VARSTR_HEADER_SIZE] = {0};
S
Shengliang Guan 已提交
946 947 948
    char trigger2[20] = {0};
    mndShowStreamTrigger(trigger2, pStream);
    STR_WITH_MAXSIZE_TO_VARSTR(trigger, trigger2, sizeof(trigger));
L
Liu Jicong 已提交
949
    pColInfo = taosArrayGet(pBlock->pDataBlock, cols++);
L
Liu Jicong 已提交
950
    colDataAppend(pColInfo, numOfRows, (const char *)&trigger, false);
L
Liu Jicong 已提交
951 952 953

    numOfRows++;
    sdbRelease(pSdb, pStream);
L
Liu Jicong 已提交
954
  }
L
Liu Jicong 已提交
955 956 957

  pShow->numOfRows += numOfRows;
  return numOfRows;
L
Liu Jicong 已提交
958 959 960 961 962 963
}

static void mndCancelGetNextStream(SMnode *pMnode, void *pIter) {
  SSdb *pSdb = pMnode->pSdb;
  sdbCancelFetch(pSdb, pIter);
}
964 965 966 967 968 969 970 971 972 973 974 975 976 977 978 979 980 981 982 983 984 985 986 987 988 989 990 991 992 993 994 995 996 997 998 999 1000 1001 1002 1003 1004 1005 1006 1007 1008 1009 1010 1011 1012 1013 1014 1015 1016 1017 1018 1019 1020 1021 1022 1023 1024 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

static int32_t mndRetrieveStreamTask(SRpcMsg *pReq, SShowObj *pShow, SSDataBlock *pBlock, int32_t rowsCapacity) {
  SMnode     *pMnode = pReq->info.node;
  SSdb       *pSdb = pMnode->pSdb;
  int32_t     numOfRows = 0;
  SStreamObj *pStream = NULL;

  while (numOfRows < rowsCapacity) {
    pShow->pIter = sdbFetch(pSdb, SDB_STREAM, pShow->pIter, (void **)&pStream);
    if (pShow->pIter == NULL) break;

    // lock
    taosRLockLatch(&pStream->lock);
    // count task num
    int32_t sz = taosArrayGetSize(pStream->tasks);
    int32_t count = 0;
    for (int32_t i = 0; i < sz; i++) {
      SArray *pLevel = taosArrayGetP(pStream->tasks, i);
      count += taosArrayGetSize(pLevel);
    }

    if (numOfRows + count > rowsCapacity) {
      blockDataEnsureCapacity(pBlock, numOfRows + count);
    }
    // add row for each task
    for (int32_t i = 0; i < sz; i++) {
      SArray *pLevel = taosArrayGetP(pStream->tasks, i);
      int32_t levelCnt = taosArrayGetSize(pLevel);
      for (int32_t j = 0; j < levelCnt; j++) {
        SStreamTask *pTask = taosArrayGetP(pLevel, j);

        SColumnInfoData *pColInfo;
        int32_t          cols = 0;

        // stream name
        char streamName[TSDB_TABLE_NAME_LEN + VARSTR_HEADER_SIZE] = {0};
        STR_WITH_MAXSIZE_TO_VARSTR(streamName, mndGetDbStr(pStream->name), sizeof(streamName));
        pColInfo = taosArrayGet(pBlock->pDataBlock, cols++);
        colDataAppend(pColInfo, numOfRows, (const char *)streamName, false);

        // task id
        pColInfo = taosArrayGet(pBlock->pDataBlock, cols++);
        colDataAppend(pColInfo, numOfRows, (const char *)&pTask->taskId, false);

        // node type
        char nodeType[20 + VARSTR_HEADER_SIZE] = {0};
        varDataSetLen(nodeType, 5);
        pColInfo = taosArrayGet(pBlock->pDataBlock, cols++);
        if (pTask->nodeId > 0) {
          memcpy(varDataVal(nodeType), "vnode", 5);
        } else {
          memcpy(varDataVal(nodeType), "snode", 5);
        }
        colDataAppend(pColInfo, numOfRows, nodeType, false);

        // node id
        pColInfo = taosArrayGet(pBlock->pDataBlock, cols++);
        int32_t nodeId = TMAX(pTask->nodeId, 0);
        colDataAppend(pColInfo, numOfRows, (const char *)&nodeId, false);

        // level
        char level[20 + VARSTR_HEADER_SIZE] = {0};
        if (pTask->taskLevel == TASK_LEVEL__SOURCE) {
          memcpy(varDataVal(level), "source", 6);
          varDataSetLen(level, 6);
        } else if (pTask->taskLevel == TASK_LEVEL__AGG) {
          memcpy(varDataVal(level), "agg", 3);
          varDataSetLen(level, 3);
        } else if (pTask->taskLevel == TASK_LEVEL__SINK) {
          memcpy(varDataVal(level), "sink", 4);
          varDataSetLen(level, 4);
        } else if (pTask->taskLevel == TASK_LEVEL__SINK) {
        }
        pColInfo = taosArrayGet(pBlock->pDataBlock, cols++);
        colDataAppend(pColInfo, numOfRows, (const char *)&level, false);

        // status
        char status[20 + VARSTR_HEADER_SIZE] = {0};
        char status2[20] = {0};
        strcpy(status, "normal");
        STR_WITH_MAXSIZE_TO_VARSTR(status, status2, sizeof(status));
        pColInfo = taosArrayGet(pBlock->pDataBlock, cols++);
        colDataAppend(pColInfo, numOfRows, (const char *)&status, false);

        numOfRows++;
      }
    }

    // unlock
    taosRUnLockLatch(&pStream->lock);

    sdbRelease(pSdb, pStream);
  }

  pShow->numOfRows += numOfRows;
  return numOfRows;
}

static void mndCancelGetNextStreamTask(SMnode *pMnode, void *pIter) {
  SSdb *pSdb = pMnode->pSdb;
  sdbCancelFetch(pSdb, pIter);
}