tq.c 42.5 KB
Newer Older
H
refact  
Hongze Cheng 已提交
1 2 3 4 5 6 7 8 9 10 11 12 13
/*
 * Copyright (c) 2019 TAOS Data, Inc. <jhtao@taosdata.com>
 *
 * This program is free software: you can use, redistribute, and/or modify
 * it under the terms of the GNU Affero General Public License, version 3
 * or later ("AGPL"), as published by the Free Software Foundation.
 *
 * This program is distributed in the hope that it will be useful, but WITHOUT
 * ANY WARRANTY; without even the implied warranty of MERCHANTABILITY or
 * FITNESS FOR A PARTICULAR PURPOSE.
 *
 * You should have received a copy of the GNU Affero General Public License
 * along with this program. If not, see <http://www.gnu.org/licenses/>.
S
Shengliang Guan 已提交
14 15
 */

H
Hongze Cheng 已提交
16
#include "tq.h"
S
Shengliang Guan 已提交
17

L
Liu Jicong 已提交
18
int32_t tqInit() {
L
Liu Jicong 已提交
19 20 21 22 23 24
  int8_t old;
  while (1) {
    old = atomic_val_compare_exchange_8(&tqMgmt.inited, 0, 2);
    if (old != 2) break;
  }

25 26 27 28 29 30
  if (old == 0) {
    tqMgmt.timer = taosTmrInit(10000, 100, 10000, "TQ");
    if (tqMgmt.timer == NULL) {
      atomic_store_8(&tqMgmt.inited, 0);
      return -1;
    }
31 32 33
    if (streamInit() < 0) {
      return -1;
    }
L
Liu Jicong 已提交
34
    atomic_store_8(&tqMgmt.inited, 1);
35
  }
36

L
Liu Jicong 已提交
37 38
  return 0;
}
L
Liu Jicong 已提交
39

40
void tqCleanUp() {
L
Liu Jicong 已提交
41 42 43 44 45 46 47 48
  int8_t old;
  while (1) {
    old = atomic_val_compare_exchange_8(&tqMgmt.inited, 1, 2);
    if (old != 2) break;
  }

  if (old == 1) {
    taosTmrCleanUp(tqMgmt.timer);
L
Liu Jicong 已提交
49
    streamCleanUp();
L
Liu Jicong 已提交
50 51
    atomic_store_8(&tqMgmt.inited, 0);
  }
52
}
L
Liu Jicong 已提交
53

54 55 56 57 58 59 60 61
static void destroySTqHandle(void* data) {
  STqHandle* pData = (STqHandle*)data;
  qDestroyTask(pData->execHandle.task);
  if (pData->execHandle.subType == TOPIC_SUB_TYPE__COLUMN) {
  } else if (pData->execHandle.subType == TOPIC_SUB_TYPE__DB) {
    tqCloseReader(pData->execHandle.pExecReader);
    walCloseReader(pData->pWalReader);
    taosHashCleanup(pData->execHandle.execDb.pFilterOutTbUid);
L
Liu Jicong 已提交
62
  } else if (pData->execHandle.subType == TOPIC_SUB_TYPE__TABLE) {
63 64 65 66 67
    walCloseReader(pData->pWalReader);
    tqCloseReader(pData->execHandle.pExecReader);
  }
}

L
Liu Jicong 已提交
68
static void tqPushEntryFree(void* data) {
L
Liu Jicong 已提交
69
  STqPushEntry* p = *(void**)data;
L
Liu Jicong 已提交
70 71 72
  taosMemoryFree(p);
}

L
Liu Jicong 已提交
73
STQ* tqOpen(const char* path, SVnode* pVnode) {
74
  STQ* pTq = taosMemoryCalloc(1, sizeof(STQ));
L
Liu Jicong 已提交
75
  if (pTq == NULL) {
L
Liu Jicong 已提交
76
    terrno = TSDB_CODE_TQ_OUT_OF_MEMORY;
L
Liu Jicong 已提交
77 78
    return NULL;
  }
H
Hongze Cheng 已提交
79
  pTq->path = strdup(path);
L
Liu Jicong 已提交
80
  pTq->pVnode = pVnode;
81

82
  pTq->pHandle = taosHashInit(64, MurmurHash3_32, true, HASH_ENTRY_LOCK);
L
Liu Jicong 已提交
83

84 85
  taosHashSetFreeFp(pTq->pHandle, destroySTqHandle);

L
Liu Jicong 已提交
86 87
  taosInitRWLatch(&pTq->pushLock);
  pTq->pPushMgr = taosHashInit(64, taosGetDefaultHashFunction(TSDB_DATA_TYPE_BIGINT), true, HASH_NO_LOCK);
L
Liu Jicong 已提交
88
  taosHashSetFreeFp(pTq->pPushMgr, tqPushEntryFree);
L
Liu Jicong 已提交
89

90
  pTq->pCheckInfo = taosHashInit(64, MurmurHash3_32, true, HASH_ENTRY_LOCK);
L
Liu Jicong 已提交
91

L
Liu Jicong 已提交
92
  if (tqMetaOpen(pTq) < 0) {
93 94 95
    ASSERT(0);
  }

L
Liu Jicong 已提交
96 97
  pTq->pOffsetStore = tqOffsetOpen(pTq);
  if (pTq->pOffsetStore == NULL) {
98 99 100
    ASSERT(0);
  }

101
  pTq->pStreamMeta = streamMetaOpen(path, pTq, (FTaskExpand*)tqExpandTask, pTq->pVnode->config.vgId);
L
Liu Jicong 已提交
102 103 104 105
  if (pTq->pStreamMeta == NULL) {
    ASSERT(0);
  }

L
Liu Jicong 已提交
106 107 108 109
  if (streamLoadTasks(pTq->pStreamMeta) < 0) {
    ASSERT(0);
  }

L
Liu Jicong 已提交
110 111
  return pTq;
}
L
Liu Jicong 已提交
112

L
Liu Jicong 已提交
113
void tqClose(STQ* pTq) {
H
Hongze Cheng 已提交
114
  if (pTq) {
115
    tqOffsetClose(pTq->pOffsetStore);
116 117 118
    taosHashCleanup(pTq->pHandle);
    taosHashCleanup(pTq->pPushMgr);
    taosHashCleanup(pTq->pCheckInfo);
119
    taosMemoryFree(pTq->path);
L
Liu Jicong 已提交
120
    tqMetaClose(pTq);
L
Liu Jicong 已提交
121
    streamMetaClose(pTq->pStreamMeta);
wafwerar's avatar
wafwerar 已提交
122
    taosMemoryFree(pTq);
H
Hongze Cheng 已提交
123
  }
L
Liu Jicong 已提交
124
}
L
Liu Jicong 已提交
125

L
Liu Jicong 已提交
126
int32_t tqSendMetaPollRsp(STQ* pTq, const SRpcMsg* pMsg, const SMqPollReq* pReq, const SMqMetaRsp* pRsp) {
127 128 129 130 131 132 133
  int32_t len = 0;
  int32_t code = 0;
  tEncodeSize(tEncodeSMqMetaRsp, pRsp, len, code);
  if (code < 0) {
    return -1;
  }
  int32_t tlen = sizeof(SMqRspHead) + len;
L
Liu Jicong 已提交
134 135 136 137 138 139 140 141 142 143
  void*   buf = rpcMallocCont(tlen);
  if (buf == NULL) {
    return -1;
  }

  ((SMqRspHead*)buf)->mqMsgType = TMQ_MSG_TYPE__POLL_META_RSP;
  ((SMqRspHead*)buf)->epoch = pReq->epoch;
  ((SMqRspHead*)buf)->consumerId = pReq->consumerId;

  void* abuf = POINTER_SHIFT(buf, sizeof(SMqRspHead));
144 145 146 147 148

  SEncoder encoder = {0};
  tEncoderInit(&encoder, abuf, len);
  tEncodeSMqMetaRsp(&encoder, pRsp);
  tEncoderClear(&encoder);
L
Liu Jicong 已提交
149 150 151 152 153 154 155 156 157

  SRpcMsg resp = {
      .info = pMsg->info,
      .pCont = buf,
      .contLen = tlen,
      .code = 0,
  };
  tmsgSendRsp(&resp);

158 159
  tqDebug("vgId:%d, from consumer:%" PRId64 ", (epoch %d) send rsp, res msg type %d, offset type:%d",
          TD_VID(pTq->pVnode), pReq->consumerId, pReq->epoch, pRsp->resMsgType, pRsp->rspOffset.type);
L
Liu Jicong 已提交
160 161 162 163

  return 0;
}

L
Liu Jicong 已提交
164 165 166 167 168 169 170 171 172 173
int32_t tqPushDataRsp(STQ* pTq, STqPushEntry* pPushEntry) {
  SMqDataRsp* pRsp = &pPushEntry->dataRsp;

  ASSERT(taosArrayGetSize(pRsp->blockData) == pRsp->blockNum);
  ASSERT(taosArrayGetSize(pRsp->blockDataLen) == pRsp->blockNum);

  ASSERT(!pRsp->withSchema);
  ASSERT(taosArrayGetSize(pRsp->blockSchema) == 0);

  if (pRsp->reqOffset.type == TMQ_OFFSET__LOG) {
L
Liu Jicong 已提交
174 175 176 177 178
    /*if (pRsp->blockNum > 0) {*/
    /*ASSERT(pRsp->rspOffset.version > pRsp->reqOffset.version);*/
    /*} else {*/
    ASSERT(pRsp->rspOffset.version > pRsp->reqOffset.version);
    /*}*/
L
Liu Jicong 已提交
179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194
  }

  int32_t len = 0;
  int32_t code = 0;
  tEncodeSize(tEncodeSMqDataRsp, pRsp, len, code);

  if (code < 0) {
    return -1;
  }

  int32_t tlen = sizeof(SMqRspHead) + len;
  void*   buf = rpcMallocCont(tlen);
  if (buf == NULL) {
    return -1;
  }

L
Liu Jicong 已提交
195
  memcpy(buf, &pPushEntry->dataRsp.head, sizeof(SMqRspHead));
L
Liu Jicong 已提交
196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 212 213 214 215 216 217

  void* abuf = POINTER_SHIFT(buf, sizeof(SMqRspHead));

  SEncoder encoder = {0};
  tEncoderInit(&encoder, abuf, len);
  tEncodeSMqDataRsp(&encoder, pRsp);
  tEncoderClear(&encoder);

  SRpcMsg rsp = {
      .info = pPushEntry->pInfo,
      .pCont = buf,
      .contLen = tlen,
      .code = 0,
  };

  tmsgSendRsp(&rsp);

  char buf1[80] = {0};
  char buf2[80] = {0};
  tFormatOffset(buf1, 80, &pRsp->reqOffset);
  tFormatOffset(buf2, 80, &pRsp->rspOffset);
  tqDebug("vgId:%d, from consumer:%" PRId64 ", (epoch %d) push rsp, block num: %d, reqOffset:%s, rspOffset:%s",
L
Liu Jicong 已提交
218
          TD_VID(pTq->pVnode), pRsp->head.consumerId, pRsp->head.epoch, pRsp->blockNum, buf1, buf2);
L
Liu Jicong 已提交
219 220 221 222

  return 0;
}

L
Liu Jicong 已提交
223 224 225 226
int32_t tqSendDataRsp(STQ* pTq, const SRpcMsg* pMsg, const SMqPollReq* pReq, const SMqDataRsp* pRsp) {
  ASSERT(taosArrayGetSize(pRsp->blockData) == pRsp->blockNum);
  ASSERT(taosArrayGetSize(pRsp->blockDataLen) == pRsp->blockNum);

L
Liu Jicong 已提交
227 228
  ASSERT(!pRsp->withSchema);
  ASSERT(taosArrayGetSize(pRsp->blockSchema) == 0);
L
Liu Jicong 已提交
229

230 231 232 233 234 235
  if (pRsp->reqOffset.type == TMQ_OFFSET__LOG) {
    if (pRsp->blockNum > 0) {
      ASSERT(pRsp->rspOffset.version > pRsp->reqOffset.version);
    } else {
      ASSERT(pRsp->rspOffset.version >= pRsp->reqOffset.version);
    }
L
Liu Jicong 已提交
236 237
  }

wmmhello's avatar
wmmhello 已提交
238 239
  int32_t len = 0;
  int32_t code = 0;
L
Liu Jicong 已提交
240 241 242 243 244
  tEncodeSize(tEncodeSMqDataRsp, pRsp, len, code);
  if (code < 0) {
    return -1;
  }
  int32_t tlen = sizeof(SMqRspHead) + len;
L
Liu Jicong 已提交
245 246 247 248 249 250 251 252 253 254 255
  void*   buf = rpcMallocCont(tlen);
  if (buf == NULL) {
    return -1;
  }

  ((SMqRspHead*)buf)->mqMsgType = TMQ_MSG_TYPE__POLL_RSP;
  ((SMqRspHead*)buf)->epoch = pReq->epoch;
  ((SMqRspHead*)buf)->consumerId = pReq->consumerId;

  void* abuf = POINTER_SHIFT(buf, sizeof(SMqRspHead));

wmmhello's avatar
wmmhello 已提交
256
  SEncoder encoder = {0};
L
Liu Jicong 已提交
257
  tEncoderInit(&encoder, abuf, len);
L
Liu Jicong 已提交
258
  tEncodeSMqDataRsp(&encoder, pRsp);
wmmhello's avatar
wmmhello 已提交
259
  tEncoderClear(&encoder);
L
Liu Jicong 已提交
260 261

  SRpcMsg rsp = {
L
Liu Jicong 已提交
262 263 264 265 266
      .info = pMsg->info,
      .pCont = buf,
      .contLen = tlen,
      .code = 0,
  };
L
Liu Jicong 已提交
267
  tmsgSendRsp(&rsp);
L
Liu Jicong 已提交
268

wmmhello's avatar
wmmhello 已提交
269 270
  char buf1[80] = {0};
  char buf2[80] = {0};
L
Liu Jicong 已提交
271 272
  tFormatOffset(buf1, 80, &pRsp->reqOffset);
  tFormatOffset(buf2, 80, &pRsp->rspOffset);
S
Shengliang Guan 已提交
273
  tqDebug("vgId:%d, from consumer:%" PRId64 ", (epoch %d) send rsp, block num: %d, reqOffset:%s, rspOffset:%s",
L
Liu Jicong 已提交
274
          TD_VID(pTq->pVnode), pReq->consumerId, pReq->epoch, pRsp->blockNum, buf1, buf2);
L
Liu Jicong 已提交
275 276 277 278

  return 0;
}

279 280 281 282 283 284 285 286 287 288 289 290 291 292 293 294 295 296 297 298 299 300 301 302 303 304 305 306 307 308 309 310 311 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 337 338
int32_t tqSendTaosxRsp(STQ* pTq, const SRpcMsg* pMsg, const SMqPollReq* pReq, const STaosxRsp* pRsp) {
  ASSERT(taosArrayGetSize(pRsp->blockData) == pRsp->blockNum);
  ASSERT(taosArrayGetSize(pRsp->blockDataLen) == pRsp->blockNum);

  if (pRsp->withSchema) {
    ASSERT(taosArrayGetSize(pRsp->blockSchema) == pRsp->blockNum);
  } else {
    ASSERT(taosArrayGetSize(pRsp->blockSchema) == 0);
  }

  if (pRsp->reqOffset.type == TMQ_OFFSET__LOG) {
    if (pRsp->blockNum > 0) {
      ASSERT(pRsp->rspOffset.version > pRsp->reqOffset.version);
    } else {
      ASSERT(pRsp->rspOffset.version >= pRsp->reqOffset.version);
    }
  }

  int32_t len = 0;
  int32_t code = 0;
  tEncodeSize(tEncodeSTaosxRsp, pRsp, len, code);
  if (code < 0) {
    return -1;
  }
  int32_t tlen = sizeof(SMqRspHead) + len;
  void*   buf = rpcMallocCont(tlen);
  if (buf == NULL) {
    return -1;
  }

  ((SMqRspHead*)buf)->mqMsgType = TMQ_MSG_TYPE__TAOSX_RSP;
  ((SMqRspHead*)buf)->epoch = pReq->epoch;
  ((SMqRspHead*)buf)->consumerId = pReq->consumerId;

  void* abuf = POINTER_SHIFT(buf, sizeof(SMqRspHead));

  SEncoder encoder = {0};
  tEncoderInit(&encoder, abuf, len);
  tEncodeSTaosxRsp(&encoder, pRsp);
  tEncoderClear(&encoder);

  SRpcMsg rsp = {
      .info = pMsg->info,
      .pCont = buf,
      .contLen = tlen,
      .code = 0,
  };
  tmsgSendRsp(&rsp);

  char buf1[80] = {0};
  char buf2[80] = {0};
  tFormatOffset(buf1, 80, &pRsp->reqOffset);
  tFormatOffset(buf2, 80, &pRsp->rspOffset);
  tqDebug("taosx rsp, vgId:%d, from consumer:%" PRId64
          ", (epoch %d) send rsp, block num: %d, reqOffset:%s, rspOffset:%s",
          TD_VID(pTq->pVnode), pReq->consumerId, pReq->epoch, pRsp->blockNum, buf1, buf2);

  return 0;
}

339 340 341 342 343 344
static FORCE_INLINE bool tqOffsetLessOrEqual(const STqOffset* pLeft, const STqOffset* pRight) {
  return pLeft->val.type == TMQ_OFFSET__LOG && pRight->val.type == TMQ_OFFSET__LOG &&
         pLeft->val.version <= pRight->val.version;
}

int32_t tqProcessOffsetCommitReq(STQ* pTq, int64_t version, char* msg, int32_t msgLen) {
345 346 347 348 349 350 351 352 353
  STqOffset offset = {0};
  SDecoder  decoder;
  tDecoderInit(&decoder, msg, msgLen);
  if (tDecodeSTqOffset(&decoder, &offset) < 0) {
    ASSERT(0);
    return -1;
  }
  tDecoderClear(&decoder);

wmmhello's avatar
wmmhello 已提交
354
  if (offset.val.type == TMQ_OFFSET__SNAPSHOT_DATA || offset.val.type == TMQ_OFFSET__SNAPSHOT_META) {
L
Liu Jicong 已提交
355 356
    tqDebug("receive offset commit msg to %s on vgId:%d, offset(type:snapshot) uid:%" PRId64 ", ts:%" PRId64,
            offset.subKey, TD_VID(pTq->pVnode), offset.val.uid, offset.val.ts);
L
Liu Jicong 已提交
357
  } else if (offset.val.type == TMQ_OFFSET__LOG) {
S
Shengliang Guan 已提交
358
    tqDebug("receive offset commit msg to %s on vgId:%d, offset(type:log) version:%" PRId64, offset.subKey,
L
Liu Jicong 已提交
359
            TD_VID(pTq->pVnode), offset.val.version);
360 361 362
    if (offset.val.version + 1 == version) {
      offset.val.version += 1;
    }
363 364 365
  } else {
    ASSERT(0);
  }
366 367 368 369 370
  STqOffset* pOffset = tqOffsetRead(pTq->pOffsetStore, offset.subKey);
  if (pOffset != NULL && tqOffsetLessOrEqual(&offset, pOffset)) {
    return 0;
  }

371 372 373
  if (tqOffsetWrite(pTq->pOffsetStore, &offset) < 0) {
    ASSERT(0);
    return -1;
374
  }
375 376

  if (offset.val.type == TMQ_OFFSET__LOG) {
377
    STqHandle* pHandle = taosHashGet(pTq->pHandle, offset.subKey, strlen(offset.subKey));
L
Liu Jicong 已提交
378 379 380 381 382
    if (pHandle) {
      if (walRefVer(pHandle->pRef, offset.val.version) < 0) {
        ASSERT(0);
        return -1;
      }
383 384 385
    }
  }

386 387
  // rsp

388 389
  /*}*/
  /*}*/
390 391 392 393

  return 0;
}

L
Liu Jicong 已提交
394
int32_t tqCheckColModifiable(STQ* pTq, int64_t tbUid, int32_t colId) {
L
Liu Jicong 已提交
395 396
  void* pIter = NULL;
  while (1) {
397
    pIter = taosHashIterate(pTq->pCheckInfo, pIter);
L
Liu Jicong 已提交
398
    if (pIter == NULL) break;
399
    STqCheckInfo* pCheck = (STqCheckInfo*)pIter;
L
Liu Jicong 已提交
400 401
    if (pCheck->ntbUid == tbUid) {
      int32_t sz = taosArrayGetSize(pCheck->colIdList);
L
Liu Jicong 已提交
402
      for (int32_t i = 0; i < sz; i++) {
L
Liu Jicong 已提交
403 404
        int16_t forbidColId = *(int16_t*)taosArrayGet(pCheck->colIdList, i);
        if (forbidColId == colId) {
405
          taosHashCancelIterate(pTq->pCheckInfo, pIter);
L
Liu Jicong 已提交
406 407 408 409 410 411 412 413
          return -1;
        }
      }
    }
  }
  return 0;
}

L
Liu Jicong 已提交
414 415 416 417 418 419 420 421 422 423
static int32_t tqInitDataRsp(SMqDataRsp* pRsp, const SMqPollReq* pReq, int8_t subType) {
  pRsp->reqOffset = pReq->reqOffset;

  pRsp->blockData = taosArrayInit(0, sizeof(void*));
  pRsp->blockDataLen = taosArrayInit(0, sizeof(int32_t));

  if (pRsp->blockData == NULL || pRsp->blockDataLen == NULL) {
    return -1;
  }

L
Liu Jicong 已提交
424 425
  pRsp->withTbName = 0;
#if 0
L
Liu Jicong 已提交
426 427 428 429 430 431 432 433
  pRsp->withTbName = pReq->withTbName;
  if (pRsp->withTbName) {
    pRsp->blockTbName = taosArrayInit(0, sizeof(void*));
    if (pRsp->blockTbName == NULL) {
      // TODO free
      return -1;
    }
  }
L
Liu Jicong 已提交
434
#endif
L
Liu Jicong 已提交
435

L
Liu Jicong 已提交
436 437 438
  ASSERT(subType == TOPIC_SUB_TYPE__COLUMN);
  pRsp->withSchema = false;

L
Liu Jicong 已提交
439 440 441
  return 0;
}

442 443 444 445 446 447 448 449 450 451 452 453 454 455 456 457
static int32_t tqInitTaosxRsp(STaosxRsp* pRsp, const SMqPollReq* pReq) {
  pRsp->reqOffset = pReq->reqOffset;

  pRsp->withTbName = 1;
  pRsp->withSchema = 1;
  pRsp->blockData = taosArrayInit(0, sizeof(void*));
  pRsp->blockDataLen = taosArrayInit(0, sizeof(int32_t));
  pRsp->blockTbName = taosArrayInit(0, sizeof(void*));
  pRsp->blockSchema = taosArrayInit(0, sizeof(void*));

  if (pRsp->blockData == NULL || pRsp->blockDataLen == NULL || pRsp->blockTbName == NULL || pRsp->blockSchema == NULL) {
    return -1;
  }
  return 0;
}

L
Liu Jicong 已提交
458
int32_t tqProcessPollReq(STQ* pTq, SRpcMsg* pMsg) {
L
Liu Jicong 已提交
459 460 461 462 463 464
  SMqPollReq*  pReq = pMsg->pCont;
  int64_t      consumerId = pReq->consumerId;
  int32_t      reqEpoch = pReq->epoch;
  int32_t      code = 0;
  STqOffsetVal reqOffset = pReq->reqOffset;
  STqOffsetVal fetchOffsetNew;
wmmhello's avatar
wmmhello 已提交
465
  SWalCkHead*  pCkHead = NULL;
L
Liu Jicong 已提交
466 467

  // 1.find handle
468
  STqHandle* pHandle = taosHashGet(pTq->pHandle, pReq->subKey, strlen(pReq->subKey));
L
Liu Jicong 已提交
469 470
  /*ASSERT(pHandle);*/
  if (pHandle == NULL) {
S
Shengliang Guan 已提交
471 472
    tqError("tmq poll: no consumer handle for consumer:%" PRId64 ", in vgId:%d, subkey %s", consumerId,
            TD_VID(pTq->pVnode), pReq->subKey);
L
Liu Jicong 已提交
473 474 475 476 477
    return -1;
  }

  // check rebalance
  if (pHandle->consumerId != consumerId) {
S
Shengliang Guan 已提交
478 479
    tqError("tmq poll: consumer handle mismatch for consumer:%" PRId64
            ", in vgId:%d, subkey %s, handle consumer id %" PRId64,
L
Liu Jicong 已提交
480
            consumerId, TD_VID(pTq->pVnode), pReq->subKey, pHandle->consumerId);
L
Liu Jicong 已提交
481
    terrno = TSDB_CODE_TMQ_CONSUMER_MISMATCH;
L
Liu Jicong 已提交
482 483 484 485 486 487 488 489 490
    return -1;
  }

  // update epoch if need
  int32_t consumerEpoch = atomic_load_32(&pHandle->epoch);
  while (consumerEpoch < reqEpoch) {
    consumerEpoch = atomic_val_compare_exchange_32(&pHandle->epoch, consumerEpoch, reqEpoch);
  }

L
Liu Jicong 已提交
491 492
  char buf[80];
  tFormatOffset(buf, 80, &reqOffset);
S
Shengliang Guan 已提交
493
  tqDebug("tmq poll: consumer %" PRId64 " (epoch %d), subkey %s, recv poll req in vg %d, req offset %s", consumerId,
L
Liu Jicong 已提交
494 495
          pReq->epoch, pHandle->subKey, TD_VID(pTq->pVnode), buf);

L
Liu Jicong 已提交
496 497 498 499 500 501 502
  // 2.reset offset if needed
  if (reqOffset.type > 0) {
    fetchOffsetNew = reqOffset;
  } else {
    STqOffset* pOffset = tqOffsetRead(pTq->pOffsetStore, pReq->subKey);
    if (pOffset != NULL) {
      fetchOffsetNew = pOffset->val;
L
Liu Jicong 已提交
503 504
      char formatBuf[80];
      tFormatOffset(formatBuf, 80, &fetchOffsetNew);
L
Liu Jicong 已提交
505 506
      tqDebug("tmq poll: consumer %" PRId64 ", subkey %s, vg %d, offset reset to %s", consumerId, pHandle->subKey,
              TD_VID(pTq->pVnode), formatBuf);
L
Liu Jicong 已提交
507 508
    } else {
      if (reqOffset.type == TMQ_OFFSET__RESET_EARLIEAST) {
L
Liu Jicong 已提交
509 510
        if (pReq->useSnapshot) {
          if (pHandle->fetchMeta) {
wmmhello's avatar
wmmhello 已提交
511
            tqOffsetResetToMeta(&fetchOffsetNew, 0);
L
Liu Jicong 已提交
512
          } else {
513
            tqOffsetResetToData(&fetchOffsetNew, 0, 0);
L
Liu Jicong 已提交
514 515 516 517 518
          }
        } else {
          tqOffsetResetToLog(&fetchOffsetNew, walGetFirstVer(pTq->pVnode->pWal));
        }
      } else if (reqOffset.type == TMQ_OFFSET__RESET_LATEST) {
519 520 521
        SMqDataRsp dataRsp = {0};
        tqInitDataRsp(&dataRsp, pReq, pHandle->execHandle.subType);

L
Liu Jicong 已提交
522
        tqOffsetResetToLog(&dataRsp.rspOffset, walGetLastVer(pTq->pVnode->pWal));
S
Shengliang Guan 已提交
523 524
        tqDebug("tmq poll: consumer %" PRId64 ", subkey %s, vg %d, offset reset to %" PRId64, consumerId,
                pHandle->subKey, TD_VID(pTq->pVnode), dataRsp.rspOffset.version);
L
Liu Jicong 已提交
525 526 527
        if (tqSendDataRsp(pTq, pMsg, pReq, &dataRsp) < 0) {
          code = -1;
        }
L
Liu Jicong 已提交
528 529
        tDeleteSMqDataRsp(&dataRsp);
        return code;
L
Liu Jicong 已提交
530
      } else if (reqOffset.type == TMQ_OFFSET__RESET_NONE) {
531 532
        tqError("tmq poll: subkey %s, no offset committed for consumer %" PRId64
                " in vg %d, subkey %s, reset none failed",
L
Liu Jicong 已提交
533
                pHandle->subKey, consumerId, TD_VID(pTq->pVnode), pReq->subKey);
L
Liu Jicong 已提交
534
        terrno = TSDB_CODE_TQ_NO_COMMITTED_OFFSET;
535
        return -1;
L
Liu Jicong 已提交
536 537 538 539
      }
    }
  }

540 541 542
  if (pHandle->execHandle.subType == TOPIC_SUB_TYPE__COLUMN) {
    SMqDataRsp dataRsp = {0};
    tqInitDataRsp(&dataRsp, pReq, pHandle->execHandle.subType);
L
Liu Jicong 已提交
543 544
    // lock
    taosWLockLatch(&pTq->pushLock);
545 546
    tqScanData(pTq, pHandle, &dataRsp, &fetchOffsetNew);

L
Liu Jicong 已提交
547
#if 1
L
Liu Jicong 已提交
548
    if (dataRsp.blockNum == 0 && dataRsp.reqOffset.type == TMQ_OFFSET__LOG &&
549
        dataRsp.reqOffset.version == dataRsp.rspOffset.version) {
L
Liu Jicong 已提交
550 551 552
      STqPushEntry* pPushEntry = taosMemoryCalloc(1, sizeof(STqPushEntry));
      if (pPushEntry != NULL) {
        pPushEntry->pInfo = pMsg->info;
L
Liu Jicong 已提交
553
        memcpy(pPushEntry->subKey, pHandle->subKey, TSDB_SUBSCRIBE_KEY_LEN);
L
Liu Jicong 已提交
554
        dataRsp.withTbName = 0;
L
Liu Jicong 已提交
555
        memcpy(&pPushEntry->dataRsp, &dataRsp, sizeof(SMqDataRsp));
L
Liu Jicong 已提交
556 557 558
        pPushEntry->dataRsp.head.consumerId = consumerId;
        pPushEntry->dataRsp.head.epoch = reqEpoch;
        pPushEntry->dataRsp.head.mqMsgType = TMQ_MSG_TYPE__POLL_RSP;
L
Liu Jicong 已提交
559
        taosHashPut(pTq->pPushMgr, pHandle->subKey, strlen(pHandle->subKey) + 1, &pPushEntry, sizeof(void*));
560
        tqDebug("tmq poll: consumer %" PRId64 ", subkey %s, vg %d save handle to push mgr", consumerId, pHandle->subKey,
L
Liu Jicong 已提交
561 562 563 564 565 566 567
                TD_VID(pTq->pVnode));
        // unlock
        taosWUnLockLatch(&pTq->pushLock);
        return 0;
      }
    }
    taosWUnLockLatch(&pTq->pushLock);
L
Liu Jicong 已提交
568
#endif
L
Liu Jicong 已提交
569

570 571 572 573
    if (tqSendDataRsp(pTq, pMsg, pReq, &dataRsp) < 0) {
      code = -1;
    }

574 575
    tqDebug("tmq poll: consumer %" PRId64 ", subkey %s, vg %d, send data blockNum:%d, offset type:%d, uid:%" PRId64
            ", version:%" PRId64 "",
576 577 578 579 580 581 582 583 584 585 586 587 588 589 590 591
            consumerId, pHandle->subKey, TD_VID(pTq->pVnode), dataRsp.blockNum, dataRsp.rspOffset.type,
            dataRsp.rspOffset.uid, dataRsp.rspOffset.version);

    tDeleteSMqDataRsp(&dataRsp);
    return code;
  }

  // for taosx
  ASSERT(pHandle->execHandle.subType != TOPIC_SUB_TYPE__COLUMN);

  SMqMetaRsp metaRsp = {0};

  STaosxRsp taosxRsp = {0};
  tqInitTaosxRsp(&taosxRsp, pReq);

  if (fetchOffsetNew.type != TMQ_OFFSET__LOG) {
L
Liu Jicong 已提交
592
    tqScanTaosx(pTq, pHandle, &taosxRsp, &metaRsp, &fetchOffsetNew);
L
Liu Jicong 已提交
593

L
Liu Jicong 已提交
594
    if (metaRsp.metaRspLen > 0) {
wmmhello's avatar
wmmhello 已提交
595 596 597
      if (tqSendMetaPollRsp(pTq, pMsg, pReq, &metaRsp) < 0) {
        code = -1;
      }
598 599 600
      tqDebug("tmq poll: consumer %" PRId64 ", subkey %s, vg %d, send meta offset type:%d,uid:%" PRId64
              ",version:%" PRId64 "",
              consumerId, pHandle->subKey, TD_VID(pTq->pVnode), metaRsp.rspOffset.type, metaRsp.rspOffset.uid,
L
Liu Jicong 已提交
601
              metaRsp.rspOffset.version);
wmmhello's avatar
wmmhello 已提交
602
      taosMemoryFree(metaRsp.metaRsp);
603 604
      tDeleteSTaosxRsp(&taosxRsp);
      return code;
wmmhello's avatar
wmmhello 已提交
605 606
    }

607 608
    if (taosxRsp.blockNum > 0) {
      if (tqSendTaosxRsp(pTq, pMsg, pReq, &taosxRsp) < 0) {
wmmhello's avatar
wmmhello 已提交
609 610
        code = -1;
      }
611 612
      tDeleteSTaosxRsp(&taosxRsp);
      return code;
L
Liu Jicong 已提交
613
    } else {
614
      fetchOffsetNew = taosxRsp.rspOffset;
615
    }
wmmhello's avatar
wmmhello 已提交
616

617 618
    tqDebug("taosx poll: consumer %" PRId64 ", subkey %s, vg %d, send data blockNum:%d, offset type:%d,uid:%" PRId64
            ",version:%" PRId64 "",
619 620
            consumerId, pHandle->subKey, TD_VID(pTq->pVnode), taosxRsp.blockNum, taosxRsp.rspOffset.type,
            taosxRsp.rspOffset.uid, taosxRsp.rspOffset.version);
621 622
  }

623
  if (fetchOffsetNew.type == TMQ_OFFSET__LOG) {
wmmhello's avatar
wmmhello 已提交
624 625 626
    int64_t fetchVer = fetchOffsetNew.version + 1;
    pCkHead = taosMemoryMalloc(sizeof(SWalCkHead) + 2048);
    if (pCkHead == NULL) {
627
      tDeleteSTaosxRsp(&taosxRsp);
628
      return -1;
629 630
    }

wmmhello's avatar
wmmhello 已提交
631 632 633 634 635 636
    walSetReaderCapacity(pHandle->pWalReader, 2048);

    while (1) {
      consumerEpoch = atomic_load_32(&pHandle->epoch);
      if (consumerEpoch > reqEpoch) {
        tqWarn("tmq poll: consumer %" PRId64 " (epoch %d), subkey %s, vg %d offset %" PRId64
L
Liu Jicong 已提交
637
               ", found new consumer epoch %d, discard req epoch %d",
wmmhello's avatar
wmmhello 已提交
638 639 640 641 642
               consumerId, pReq->epoch, pHandle->subKey, TD_VID(pTq->pVnode), fetchVer, consumerEpoch, reqEpoch);
        break;
      }

      if (tqFetchLog(pTq, pHandle, &fetchVer, &pCkHead) < 0) {
643 644
        tqOffsetResetToLog(&taosxRsp.rspOffset, fetchVer);
        if (tqSendTaosxRsp(pTq, pMsg, pReq, &taosxRsp) < 0) {
wmmhello's avatar
wmmhello 已提交
645 646
          code = -1;
        }
647
        tDeleteSTaosxRsp(&taosxRsp);
L
Liu Jicong 已提交
648
        taosMemoryFreeClear(pCkHead);
649
        return code;
wmmhello's avatar
wmmhello 已提交
650 651 652
      }

      SWalCont* pHead = &pCkHead->head;
wmmhello's avatar
wmmhello 已提交
653

wmmhello's avatar
wmmhello 已提交
654 655 656 657 658 659
      tqDebug("tmq poll: consumer:%" PRId64 ", (epoch %d) iter log, vgId:%d offset %" PRId64 " msgType %d", consumerId,
              pReq->epoch, TD_VID(pTq->pVnode), fetchVer, pHead->msgType);

      if (pHead->msgType == TDMT_VND_SUBMIT) {
        SSubmitReq* pCont = (SSubmitReq*)&pHead->body;

660
        if (tqTaosxScanLog(pTq, pHandle, pCont, &taosxRsp) < 0) {
wmmhello's avatar
wmmhello 已提交
661 662 663 664
          /*ASSERT(0);*/
        }
        // TODO batch optimization:
        // TODO continue scan until meeting batch requirement
665 666 667
        if (taosxRsp.blockNum > 0 /* threshold */) {
          tqOffsetResetToLog(&taosxRsp.rspOffset, fetchVer);
          if (tqSendTaosxRsp(pTq, pMsg, pReq, &taosxRsp) < 0) {
wmmhello's avatar
wmmhello 已提交
668 669
            code = -1;
          }
670
          tDeleteSTaosxRsp(&taosxRsp);
L
Liu Jicong 已提交
671
          taosMemoryFreeClear(pCkHead);
672
          return code;
wmmhello's avatar
wmmhello 已提交
673 674 675 676 677 678 679 680 681 682 683 684 685 686
        } else {
          fetchVer++;
        }

      } else {
        ASSERT(pHandle->fetchMeta);
        ASSERT(IS_META_MSG(pHead->msgType));
        tqDebug("fetch meta msg, ver:%" PRId64 ", type:%d", pHead->version, pHead->msgType);
        tqOffsetResetToLog(&metaRsp.rspOffset, fetchVer);
        metaRsp.resMsgType = pHead->msgType;
        metaRsp.metaRspLen = pHead->bodyLen;
        metaRsp.metaRsp = pHead->body;
        if (tqSendMetaPollRsp(pTq, pMsg, pReq, &metaRsp) < 0) {
          code = -1;
L
Liu Jicong 已提交
687
          taosMemoryFreeClear(pCkHead);
688
          tDeleteSTaosxRsp(&taosxRsp);
689
          return code;
wmmhello's avatar
wmmhello 已提交
690 691
        }
        code = 0;
L
Liu Jicong 已提交
692
        taosMemoryFreeClear(pCkHead);
693
        tDeleteSTaosxRsp(&taosxRsp);
694
        return code;
wmmhello's avatar
wmmhello 已提交
695 696 697
      }
    }
  }
698
  tDeleteSTaosxRsp(&taosxRsp);
L
Liu Jicong 已提交
699
  taosMemoryFreeClear(pCkHead);
700
  return 0;
L
Liu Jicong 已提交
701 702
}

703
int32_t tqProcessVgDeleteReq(STQ* pTq, int64_t version, char* msg, int32_t msgLen) {
L
Liu Jicong 已提交
704
  SMqVDeleteReq* pReq = (SMqVDeleteReq*)msg;
L
Liu Jicong 已提交
705

L
Liu Jicong 已提交
706 707 708 709 710 711 712 713
  taosWLockLatch(&pTq->pushLock);
  int32_t code = taosHashRemove(pTq->pPushMgr, pReq->subKey, strlen(pReq->subKey));
  if (code != 0) {
    tqDebug("vgId:%d, tq remove push handle %s", pTq->pVnode->config.vgId, pReq->subKey);
  }
  taosWUnLockLatch(&pTq->pushLock);

  code = taosHashRemove(pTq->pHandle, pReq->subKey, strlen(pReq->subKey));
L
Liu Jicong 已提交
714 715 716
  if (code != 0) {
    tqError("cannot process tq delete req %s, since no such handle", pReq->subKey);
  }
717

L
Liu Jicong 已提交
718 719 720 721
  code = tqOffsetDelete(pTq->pOffsetStore, pReq->subKey);
  if (code != 0) {
    tqError("cannot process tq delete req %s, since no such offset", pReq->subKey);
  }
L
Liu Jicong 已提交
722

L
Liu Jicong 已提交
723
  if (tqMetaDeleteHandle(pTq, pReq->subKey) < 0) {
724 725
    ASSERT(0);
  }
L
Liu Jicong 已提交
726
  return 0;
L
Liu Jicong 已提交
727 728
}

729 730 731
int32_t tqProcessAddCheckInfoReq(STQ* pTq, int64_t version, char* msg, int32_t msgLen) {
  STqCheckInfo info = {0};
  SDecoder     decoder;
L
Liu Jicong 已提交
732
  tDecoderInit(&decoder, msg, msgLen);
733
  if (tDecodeSTqCheckInfo(&decoder, &info) < 0) {
L
Liu Jicong 已提交
734 735 736 737
    terrno = TSDB_CODE_OUT_OF_MEMORY;
    return -1;
  }
  tDecoderClear(&decoder);
738 739 740 741 742
  if (taosHashPut(pTq->pCheckInfo, info.topic, strlen(info.topic), &info, sizeof(STqCheckInfo)) < 0) {
    terrno = TSDB_CODE_OUT_OF_MEMORY;
    return -1;
  }
  if (tqMetaSaveCheckInfo(pTq, info.topic, msg, msgLen) < 0) {
L
Liu Jicong 已提交
743 744 745 746 747 748
    terrno = TSDB_CODE_OUT_OF_MEMORY;
    return -1;
  }
  return 0;
}

749 750 751 752 753 754 755 756 757 758 759 760 761
int32_t tqProcessDelCheckInfoReq(STQ* pTq, int64_t version, char* msg, int32_t msgLen) {
  if (taosHashRemove(pTq->pCheckInfo, msg, strlen(msg)) < 0) {
    terrno = TSDB_CODE_OUT_OF_MEMORY;
    return -1;
  }
  if (tqMetaDeleteCheckInfo(pTq, msg) < 0) {
    terrno = TSDB_CODE_OUT_OF_MEMORY;
    return -1;
  }
  return 0;
}

int32_t tqProcessVgChangeReq(STQ* pTq, int64_t version, char* msg, int32_t msgLen) {
L
Liu Jicong 已提交
762
  SMqRebVgReq req = {0};
L
Liu Jicong 已提交
763 764
  tDecodeSMqRebVgReq(msg, &req);
  // todo lock
765
  STqHandle* pHandle = taosHashGet(pTq->pHandle, req.subKey, strlen(req.subKey));
L
Liu Jicong 已提交
766
  if (pHandle == NULL) {
L
Liu Jicong 已提交
767
    if (req.oldConsumerId != -1) {
768 769
      tqError("vgId:%d, build new consumer handle %s for consumer %" PRId64 ", but old consumerId is %" PRId64 "",
              req.vgId, req.subKey, req.newConsumerId, req.oldConsumerId);
L
Liu Jicong 已提交
770
    }
L
Liu Jicong 已提交
771
    if (req.newConsumerId == -1) {
772
      tqError("vgId:%d, tq invalid rebalance request, new consumerId %" PRId64 "", req.vgId, req.newConsumerId);
L
Liu Jicong 已提交
773 774
      return 0;
    }
L
Liu Jicong 已提交
775 776
    STqHandle tqHandle = {0};
    pHandle = &tqHandle;
L
Liu Jicong 已提交
777 778
    /*taosInitRWLatch(&pExec->lock);*/

L
Liu Jicong 已提交
779 780 781
    memcpy(pHandle->subKey, req.subKey, TSDB_SUBSCRIBE_KEY_LEN);
    pHandle->consumerId = req.newConsumerId;
    pHandle->epoch = -1;
L
Liu Jicong 已提交
782

L
Liu Jicong 已提交
783
    pHandle->execHandle.subType = req.subType;
L
Liu Jicong 已提交
784
    pHandle->fetchMeta = req.withMeta;
wmmhello's avatar
wmmhello 已提交
785

786 787 788 789
    // TODO version should be assigned and refed during preprocess
    SWalRef* pRef = walRefCommittedVer(pTq->pVnode->pWal);
    if (pRef == NULL) {
      ASSERT(0);
L
Liu Jicong 已提交
790
      return -1;
791 792 793
    }
    int64_t ver = pRef->refVer;
    pHandle->pRef = pRef;
L
Liu Jicong 已提交
794

795 796 797 798 799 800 801
    SReadHandle handle = {
        .meta = pTq->pVnode->pMeta,
        .vnode = pTq->pVnode,
        .initTableReader = true,
        .initTqReader = true,
        .version = ver,
    };
wmmhello's avatar
wmmhello 已提交
802
    pHandle->snapshotVer = ver;
803

L
Liu Jicong 已提交
804
    if (pHandle->execHandle.subType == TOPIC_SUB_TYPE__COLUMN) {
L
Liu Jicong 已提交
805
      pHandle->execHandle.execCol.qmsg = req.qmsg;
L
Liu Jicong 已提交
806
      req.qmsg = NULL;
807 808

      pHandle->execHandle.task =
wmmhello's avatar
wmmhello 已提交
809
          qCreateQueueExecTaskInfo(pHandle->execHandle.execCol.qmsg, &handle, &pHandle->execHandle.numOfCols,
L
Liu Jicong 已提交
810
                                   &pHandle->execHandle.pSchemaWrapper);
811
      ASSERT(pHandle->execHandle.task);
L
Liu Jicong 已提交
812
      void* scanner = NULL;
813
      qExtractStreamScanner(pHandle->execHandle.task, &scanner);
L
Liu Jicong 已提交
814 815 816
      ASSERT(scanner);
      pHandle->execHandle.pExecReader = qExtractReaderFromStreamScanner(scanner);
      ASSERT(pHandle->execHandle.pExecReader);
L
Liu Jicong 已提交
817
    } else if (pHandle->execHandle.subType == TOPIC_SUB_TYPE__DB) {
wmmhello's avatar
wmmhello 已提交
818
      pHandle->pWalReader = walOpenReader(pTq->pVnode->pWal, NULL);
L
Liu Jicong 已提交
819
      pHandle->execHandle.pExecReader = tqOpenReader(pTq->pVnode);
L
Liu Jicong 已提交
820
      pHandle->execHandle.execDb.pFilterOutTbUid =
L
Liu Jicong 已提交
821
          taosHashInit(64, taosGetDefaultHashFunction(TSDB_DATA_TYPE_BIGINT), false, HASH_NO_LOCK);
822 823
      buildSnapContext(handle.meta, handle.version, 0, pHandle->execHandle.subType, pHandle->fetchMeta,
                       (SSnapContext**)(&handle.sContext));
824

825
      pHandle->execHandle.task = qCreateQueueExecTaskInfo(NULL, &handle, NULL, NULL);
L
Liu Jicong 已提交
826
    } else if (pHandle->execHandle.subType == TOPIC_SUB_TYPE__TABLE) {
wmmhello's avatar
wmmhello 已提交
827 828 829 830
      pHandle->pWalReader = walOpenReader(pTq->pVnode->pWal, NULL);

      pHandle->execHandle.execTb.suid = req.suid;

L
Liu Jicong 已提交
831
      SArray* tbUidList = taosArrayInit(0, sizeof(int64_t));
H
Hongze Cheng 已提交
832
      vnodeGetCtbIdList(pTq->pVnode, req.suid, tbUidList);
833
      tqDebug("vgId:%d, tq try to get all ctb, suid:%" PRId64, pTq->pVnode->config.vgId, req.suid);
L
Liu Jicong 已提交
834 835
      for (int32_t i = 0; i < taosArrayGetSize(tbUidList); i++) {
        int64_t tbUid = *(int64_t*)taosArrayGet(tbUidList, i);
S
Shengliang Guan 已提交
836
        tqDebug("vgId:%d, idx %d, uid:%" PRId64, TD_VID(pTq->pVnode), i, tbUid);
L
Liu Jicong 已提交
837
      }
L
Liu Jicong 已提交
838 839
      pHandle->execHandle.pExecReader = tqOpenReader(pTq->pVnode);
      tqReaderSetTbUidList(pHandle->execHandle.pExecReader, tbUidList);
L
Liu Jicong 已提交
840
      taosArrayDestroy(tbUidList);
wmmhello's avatar
wmmhello 已提交
841

L
Liu Jicong 已提交
842 843 844
      buildSnapContext(handle.meta, handle.version, req.suid, pHandle->execHandle.subType, pHandle->fetchMeta,
                       (SSnapContext**)(&handle.sContext));
      pHandle->execHandle.task = qCreateQueueExecTaskInfo(NULL, &handle, NULL, NULL);
L
Liu Jicong 已提交
845
    }
846
    taosHashPut(pTq->pHandle, req.subKey, strlen(req.subKey), pHandle, sizeof(STqHandle));
S
Shengliang Guan 已提交
847
    tqDebug("try to persist handle %s consumer %" PRId64, req.subKey, pHandle->consumerId);
L
Liu Jicong 已提交
848 849
    if (tqMetaSaveHandle(pTq, req.subKey, pHandle) < 0) {
      // TODO
P
plum-lihui 已提交
850
      ASSERT(0);
L
Liu Jicong 已提交
851
    }
L
Liu Jicong 已提交
852
  } else {
L
Liu Jicong 已提交
853
    /*ASSERT(pExec->consumerId == req.oldConsumerId);*/
L
Liu Jicong 已提交
854
    // TODO handle qmsg and exec modification
L
Liu Jicong 已提交
855 856 857
    atomic_store_32(&pHandle->epoch, -1);
    atomic_store_64(&pHandle->consumerId, req.newConsumerId);
    atomic_add_fetch_32(&pHandle->epoch, 1);
L
Liu Jicong 已提交
858 859
    if (tqMetaSaveHandle(pTq, req.subKey, pHandle) < 0) {
      // TODO
L
Liu Jicong 已提交
860
      ASSERT(0);
L
Liu Jicong 已提交
861
    }
L
Liu Jicong 已提交
862
    // close handle
L
Liu Jicong 已提交
863
  }
L
Liu Jicong 已提交
864

L
Liu Jicong 已提交
865
  return 0;
L
Liu Jicong 已提交
866
}
867

868
int32_t tqExpandTask(STQ* pTq, SStreamTask* pTask, int64_t ver) {
869
  if (pTask->taskLevel == TASK_LEVEL__AGG) {
L
Liu Jicong 已提交
870 871
    ASSERT(taosArrayGetSize(pTask->childEpInfo) != 0);
  }
L
Liu Jicong 已提交
872

L
Liu Jicong 已提交
873
  pTask->schedStatus = TASK_SCHED_STATUS__INACTIVE;
L
Liu Jicong 已提交
874 875 876

  pTask->inputQueue = streamQueueOpen();
  pTask->outputQueue = streamQueueOpen();
L
Liu Jicong 已提交
877 878

  if (pTask->inputQueue == NULL || pTask->outputQueue == NULL) {
L
Liu Jicong 已提交
879
    return -1;
L
Liu Jicong 已提交
880 881
  }

L
Liu Jicong 已提交
882 883 884
  pTask->inputStatus = TASK_INPUT_STATUS__NORMAL;
  pTask->outputStatus = TASK_OUTPUT_STATUS__NORMAL;

885 886
  pTask->pMsgCb = &pTq->pVnode->msgCb;

887 888
  pTask->startVer = ver;

889 890
  // expand executor
  if (pTask->taskLevel == TASK_LEVEL__SOURCE) {
891
    pTask->pState = streamStateOpen(pTq->pStreamMeta->path, pTask, false, -1, -1);
892 893 894 895
    if (pTask->pState == NULL) {
      return -1;
    }

896 897 898 899
    SReadHandle handle = {
        .meta = pTq->pVnode->pMeta,
        .vnode = pTq->pVnode,
        .initTqReader = 1,
900
        .pStateBackend = pTask->pState,
901 902 903
    };
    pTask->exec.executor = qCreateStreamExecTaskInfo(pTask->exec.qmsg, &handle);
    ASSERT(pTask->exec.executor);
904 905 906 907

    if (pTask->fillHistory) {
      pTask->taskStatus = TASK_STATUS__RECOVER_PREPARE;
    }
908
  } else if (pTask->taskLevel == TASK_LEVEL__AGG) {
909
    pTask->pState = streamStateOpen(pTq->pStreamMeta->path, pTask, false, -1, -1);
910 911 912
    if (pTask->pState == NULL) {
      return -1;
    }
913 914 915
    SReadHandle mgHandle = {
        .vnode = NULL,
        .numOfVgroups = (int32_t)taosArrayGetSize(pTask->childEpInfo),
916
        .pStateBackend = pTask->pState,
917 918
    };
    pTask->exec.executor = qCreateStreamExecTaskInfo(pTask->exec.qmsg, &mgHandle);
L
Liu Jicong 已提交
919
    ASSERT(pTask->exec.executor);
L
Liu Jicong 已提交
920
  }
L
Liu Jicong 已提交
921 922

  // sink
L
Liu Jicong 已提交
923
  /*pTask->ahandle = pTq->pVnode;*/
924
  if (pTask->outputType == TASK_OUTPUT__SMA) {
L
Liu Jicong 已提交
925
    pTask->smaSink.vnode = pTq->pVnode;
L
Liu Jicong 已提交
926
    pTask->smaSink.smaSink = smaHandleRes;
927
  } else if (pTask->outputType == TASK_OUTPUT__TABLE) {
L
Liu Jicong 已提交
928
    pTask->tbSink.vnode = pTq->pVnode;
L
Liu Jicong 已提交
929
    pTask->tbSink.tbSinkFunc = tqSinkToTablePipeline;
L
Liu Jicong 已提交
930

L
Liu Jicong 已提交
931 932
    ASSERT(pTask->tbSink.pSchemaWrapper);
    ASSERT(pTask->tbSink.pSchemaWrapper->pSchema);
L
Liu Jicong 已提交
933

L
Liu Jicong 已提交
934
    pTask->tbSink.pTSchema =
C
Cary Xu 已提交
935
        tdGetSTSChemaFromSSChema(pTask->tbSink.pSchemaWrapper->pSchema, pTask->tbSink.pSchemaWrapper->nCols, 1);
L
Liu Jicong 已提交
936
    ASSERT(pTask->tbSink.pTSchema);
L
Liu Jicong 已提交
937
  }
938 939 940

  streamSetupTrigger(pTask);

L
Liu Jicong 已提交
941
  tqInfo("expand stream task on vg %d, task id %d, child id %d", TD_VID(pTq->pVnode), pTask->taskId,
942
         pTask->selfChildId);
L
Liu Jicong 已提交
943
  return 0;
L
Liu Jicong 已提交
944
}
L
Liu Jicong 已提交
945

946
int32_t tqProcessTaskDeployReq(STQ* pTq, int64_t version, char* msg, int32_t msgLen) {
947 948 949 950 951 952 953 954 955 956 957 958 959 960 961 962 963 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 1066 1067 1068 1069 1070 1071 1072 1073 1074 1075 1076 1077 1078 1079 1080 1081 1082 1083 1084
  int32_t code;
#if 0
  code = streamMetaAddSerializedTask(pTq->pStreamMeta, version, msg, msgLen);
  if (code < 0) return code;
#endif

  // 1.deserialize msg and build task
  SStreamTask* pTask = taosMemoryCalloc(1, sizeof(SStreamTask));
  if (pTask == NULL) {
    return -1;
  }
  SDecoder decoder;
  tDecoderInit(&decoder, (uint8_t*)msg, msgLen);
  code = tDecodeSStreamTask(&decoder, pTask);
  if (code < 0) {
    tDecoderClear(&decoder);
    taosMemoryFree(pTask);
    return -1;
  }
  tDecoderClear(&decoder);

  // 2.save task
  code = streamMetaAddTask(pTq->pStreamMeta, version, pTask);
  if (code < 0) {
    return -1;
  }

  // 3.go through recover steps to fill history
  if (pTask->fillHistory) {
    streamSetParamForRecover(pTask);
    if (pTask->taskLevel == TASK_LEVEL__SOURCE) {
      streamSourceRecoverPrepareStep1(pTask, version);

      SStreamRecoverStep1Req req;
      streamBuildSourceRecover1Req(pTask, &req);

      void*   serialziedReq = (void*)&req;
      int32_t len = sizeof(SStreamRecoverStep1Req);

      SRpcMsg rpcMsg = {
          .contLen = len,
          .pCont = serialziedReq,
          .msgType = TDMT_VND_STREAM_RECOVER_STEP1,
      };

      tmsgPutToQueue(&pTq->pVnode->msgCb, STREAM_QUEUE, &rpcMsg);

    } else if (pTask->taskLevel == TASK_LEVEL__AGG) {
      streamAggRecoverPrepare(pTask);
    } else if (pTask->taskLevel == TASK_LEVEL__SINK) {
      // do nothing
    }
  }

  return 0;
}

int32_t tqProcessTaskRecover1Req(STQ* pTq, char* msg, int32_t msgLen) {
  int32_t                 code;
  SStreamRecoverStep1Req* pReq = (SStreamRecoverStep1Req*)msg;
  SStreamTask*            pTask = streamMetaGetTask(pTq->pStreamMeta, pReq->taskId);
  if (pTask == NULL) {
    return -1;
  }

  // check param
  int64_t fillVer1 = pTask->startVer;
  if (fillVer1 <= 0) {
    ASSERT(0);
    return -1;
  }

  // do recovery step 1
  streamSourceRecoverScanStep1(pTask);

  // build msg to launch next step
  SStreamRecoverStep2Req req;
  code = streamBuildSourceRecover2Req(pTask, &req);
  if (code < 0) {
    return -1;
  }

  // serialize msg
  int32_t len = sizeof(SStreamRecoverStep2Req);
  void*   serializedReq = (void*)&req;

  // dispatch msg
  SRpcMsg rpcMsg = {
      .code = 0,
      .contLen = len,
      .msgType = TDMT_VND_STREAM_RECOVER_STEP2,
      .pCont = (void*)serializedReq,
  };

  tmsgPutToQueue(&pTq->pVnode->msgCb, WRITE_QUEUE, &rpcMsg);

  return 0;
}

int32_t tqProcessTaskRecover2Req(STQ* pTq, int64_t version, char* msg, int32_t msgLen) {
  int32_t                 code;
  SStreamRecoverStep2Req* pReq = (SStreamRecoverStep2Req*)msg;
  SStreamTask*            pTask = streamMetaGetTask(pTq->pStreamMeta, pReq->taskId);
  if (pTask == NULL) {
    return -1;
  }

  // do recovery step 2
  code = streamSourceRecoverScanStep2(pTask, version);
  if (code < 0) {
    return -1;
  }

  // restore param
  code = streamRestoreParam(pTask);
  if (code < 0) {
    return -1;
  }

  // set status normal
  code = streamSetStatusNormal(pTask);
  if (code < 0) {
    return -1;
  }

  // dispatch recover finish req to all related downstream task
  code = streamDispatchRecoverFinishReq(pTask);
  if (code < 0) {
    return -1;
  }

  return 0;
}

int32_t tqProcessTaskRecoverFinishReq(STQ* pTq, char* msg, int32_t msgLen) {
  int32_t code;

  // deserialize
1085 1086 1087 1088 1089 1090 1091 1092
  int32_t                 len;
  SStreamRecoverFinishReq req;

  SDecoder decoder;
  tDecoderInit(&decoder, msg, sizeof(SStreamRecoverFinishReq));
  tDecodeSStreamRecoverFinishReq(&decoder, &req);
  tDecoderClear(&decoder);

1093
  // find task
1094 1095 1096 1097
  SStreamTask* pTask = streamMetaGetTask(pTq->pStreamMeta, req.taskId);
  if (pTask == NULL) {
    return -1;
  }
1098
  // do process request
1099 1100 1101 1102
  if (streamProcessRecoverFinishReq(pTask, req.childId) < 0) {
    return -1;
  }

1103
  return 0;
L
Liu Jicong 已提交
1104
}
L
Liu Jicong 已提交
1105

L
Liu Jicong 已提交
1106 1107 1108 1109 1110 1111 1112 1113 1114 1115 1116 1117 1118 1119 1120 1121
int32_t tqProcessDelReq(STQ* pTq, void* pReq, int32_t len, int64_t ver) {
  bool        failed = false;
  SDecoder*   pCoder = &(SDecoder){0};
  SDeleteRes* pRes = &(SDeleteRes){0};

  pRes->uidList = taosArrayInit(0, sizeof(tb_uid_t));
  if (pRes->uidList == NULL) {
    terrno = TSDB_CODE_OUT_OF_MEMORY;
    failed = true;
  }

  tDecoderInit(pCoder, pReq, len);
  tDecodeDeleteRes(pCoder, pRes);
  tDecoderClear(pCoder);

  int32_t sz = taosArrayGetSize(pRes->uidList);
L
Liu Jicong 已提交
1122
  if (sz == 0 || pRes->affectedRows == 0) {
L
Liu Jicong 已提交
1123 1124 1125 1126 1127 1128 1129 1130 1131 1132 1133 1134 1135 1136 1137 1138 1139 1140 1141 1142 1143 1144 1145 1146
    taosArrayDestroy(pRes->uidList);
    return 0;
  }
  SSDataBlock* pDelBlock = createSpecialDataBlock(STREAM_DELETE_DATA);
  blockDataEnsureCapacity(pDelBlock, sz);
  pDelBlock->info.rows = sz;
  pDelBlock->info.version = ver;

  for (int32_t i = 0; i < sz; i++) {
    // start key column
    SColumnInfoData* pStartCol = taosArrayGet(pDelBlock->pDataBlock, START_TS_COLUMN_INDEX);
    colDataAppend(pStartCol, i, (const char*)&pRes->skey, false);  // end key column
    SColumnInfoData* pEndCol = taosArrayGet(pDelBlock->pDataBlock, END_TS_COLUMN_INDEX);
    colDataAppend(pEndCol, i, (const char*)&pRes->ekey, false);
    // uid column
    SColumnInfoData* pUidCol = taosArrayGet(pDelBlock->pDataBlock, UID_COLUMN_INDEX);
    int64_t*         pUid = taosArrayGet(pRes->uidList, i);
    colDataAppend(pUidCol, i, (const char*)pUid, false);

    colDataAppendNULL(taosArrayGet(pDelBlock->pDataBlock, GROUPID_COLUMN_INDEX), i);
    colDataAppendNULL(taosArrayGet(pDelBlock->pDataBlock, CALCULATE_START_TS_COLUMN_INDEX), i);
    colDataAppendNULL(taosArrayGet(pDelBlock->pDataBlock, CALCULATE_END_TS_COLUMN_INDEX), i);
  }

L
Liu Jicong 已提交
1147 1148
  taosArrayDestroy(pRes->uidList);

L
Liu Jicong 已提交
1149 1150 1151
  int32_t* pRef = taosMemoryMalloc(sizeof(int32_t));
  *pRef = 1;

L
Liu Jicong 已提交
1152 1153 1154 1155 1156 1157 1158 1159 1160
  void* pIter = NULL;
  while (1) {
    pIter = taosHashIterate(pTq->pStreamMeta->pTasks, pIter);
    if (pIter == NULL) break;
    SStreamTask* pTask = *(SStreamTask**)pIter;
    if (pTask->taskLevel != TASK_LEVEL__SOURCE) continue;

    qDebug("delete req enqueue stream task: %d, ver: %" PRId64, pTask->taskId, ver);

L
Liu Jicong 已提交
1161 1162 1163 1164 1165 1166 1167 1168 1169
    if (!failed) {
      SStreamRefDataBlock* pRefBlock = taosAllocateQitem(sizeof(SStreamRefDataBlock), DEF_QITEM);
      pRefBlock->type = STREAM_INPUT__REF_DATA_BLOCK;
      pRefBlock->pBlock = pDelBlock;
      pRefBlock->dataRef = pRef;
      atomic_add_fetch_32(pRefBlock->dataRef, 1);

      if (streamTaskInput(pTask, (SStreamQueueItem*)pRefBlock) < 0) {
        qError("stream task input del failed, task id %d", pTask->taskId);
L
Liu Jicong 已提交
1170 1171

        taosFreeQitem(pRefBlock);
L
Liu Jicong 已提交
1172 1173
        continue;
      }
L
Liu Jicong 已提交
1174

L
Liu Jicong 已提交
1175 1176 1177 1178
      if (streamSchedExec(pTask) < 0) {
        qError("stream task launch failed, task id %d", pTask->taskId);
        continue;
      }
L
Liu Jicong 已提交
1179

L
Liu Jicong 已提交
1180 1181 1182 1183
    } else {
      streamTaskInputFail(pTask);
    }
  }
L
Liu Jicong 已提交
1184

L
Liu Jicong 已提交
1185 1186 1187 1188 1189 1190 1191 1192
  int32_t ref = atomic_sub_fetch_32(pRef, 1);
  ASSERT(ref >= 0);
  if (ref == 0) {
    taosMemoryFree(pDelBlock);
    taosMemoryFree(pRef);
  }

#if 0
L
Liu Jicong 已提交
1193 1194 1195 1196 1197 1198 1199 1200 1201 1202 1203 1204 1205 1206 1207 1208 1209 1210 1211 1212 1213 1214
    SStreamDataBlock* pStreamBlock = taosAllocateQitem(sizeof(SStreamDataBlock), DEF_QITEM);
    pStreamBlock->type = STREAM_INPUT__DATA_BLOCK;
    pStreamBlock->blocks = taosArrayInit(0, sizeof(SSDataBlock));
    SSDataBlock block = {0};
    assignOneDataBlock(&block, pDelBlock);
    block.info.type = STREAM_DELETE_DATA;
    taosArrayPush(pStreamBlock->blocks, &block);

    if (!failed) {
      if (streamTaskInput(pTask, (SStreamQueueItem*)pStreamBlock) < 0) {
        qError("stream task input del failed, task id %d", pTask->taskId);
        continue;
      }

      if (streamSchedExec(pTask) < 0) {
        qError("stream task launch failed, task id %d", pTask->taskId);
        continue;
      }
    } else {
      streamTaskInputFail(pTask);
    }
  }
L
Liu Jicong 已提交
1215
  blockDataDestroy(pDelBlock);
L
Liu Jicong 已提交
1216
#endif
L
Liu Jicong 已提交
1217 1218 1219 1220 1221

  return 0;
}

int32_t tqProcessSubmitReq(STQ* pTq, SSubmitReq* pReq, int64_t ver) {
L
Liu Jicong 已提交
1222 1223 1224
  void*              pIter = NULL;
  bool               failed = false;
  SStreamDataSubmit* pSubmit = NULL;
L
Liu Jicong 已提交
1225

L
Liu Jicong 已提交
1226
  pSubmit = streamDataSubmitNew(pReq);
L
Liu Jicong 已提交
1227
  if (pSubmit == NULL) {
L
Liu Jicong 已提交
1228 1229
    terrno = TSDB_CODE_OUT_OF_MEMORY;
    qError("failed to create data submit for stream since out of memory");
L
Liu Jicong 已提交
1230 1231 1232 1233
    failed = true;
  }

  while (1) {
L
Liu Jicong 已提交
1234
    pIter = taosHashIterate(pTq->pStreamMeta->pTasks, pIter);
L
Liu Jicong 已提交
1235
    if (pIter == NULL) break;
1236
    SStreamTask* pTask = *(SStreamTask**)pIter;
1237
    if (pTask->taskLevel != TASK_LEVEL__SOURCE) continue;
1238
    if (pTask->taskStatus == TASK_STATUS__RECOVER_PREPARE || pTask->taskStatus == TASK_STATUS__RECOVER1) continue;
L
Liu Jicong 已提交
1239

S
Shengliang Guan 已提交
1240
    qDebug("data submit enqueue stream task: %d, ver: %" PRId64, pTask->taskId, ver);
L
Liu Jicong 已提交
1241

L
Liu Jicong 已提交
1242 1243
    if (!failed) {
      if (streamTaskInput(pTask, (SStreamQueueItem*)pSubmit) < 0) {
L
Liu Jicong 已提交
1244
        qError("stream task input failed, task id %d", pTask->taskId);
L
Liu Jicong 已提交
1245 1246 1247
        continue;
      }

L
Liu Jicong 已提交
1248
      if (streamSchedExec(pTask) < 0) {
L
Liu Jicong 已提交
1249
        qError("stream task launch failed, task id %d", pTask->taskId);
L
Liu Jicong 已提交
1250 1251
        continue;
      }
L
Liu Jicong 已提交
1252
    } else {
L
Liu Jicong 已提交
1253
      streamTaskInputFail(pTask);
L
Liu Jicong 已提交
1254 1255 1256
    }
  }

L
Liu Jicong 已提交
1257
  if (pSubmit) {
L
Liu Jicong 已提交
1258
    streamDataSubmitRefDec(pSubmit);
L
Liu Jicong 已提交
1259
    taosFreeQitem(pSubmit);
L
Liu Jicong 已提交
1260
  }
L
Liu Jicong 已提交
1261 1262

  return failed ? -1 : 0;
L
Liu Jicong 已提交
1263 1264
}

L
Liu Jicong 已提交
1265
int32_t tqProcessTaskRunReq(STQ* pTq, SRpcMsg* pMsg) {
L
Liu Jicong 已提交
1266
  //
L
Liu Jicong 已提交
1267 1268
  SStreamTaskRunReq* pReq = pMsg->pCont;
  int32_t            taskId = pReq->taskId;
L
Liu Jicong 已提交
1269 1270 1271
  SStreamTask*       pTask = streamMetaGetTask(pTq->pStreamMeta, taskId);
  if (pTask) {
    streamProcessRunReq(pTask);
L
Liu Jicong 已提交
1272
    return 0;
1273 1274
  } else {
    return -1;
L
Liu Jicong 已提交
1275
  }
L
Liu Jicong 已提交
1276 1277
}

L
Liu Jicong 已提交
1278
int32_t tqProcessTaskDispatchReq(STQ* pTq, SRpcMsg* pMsg, bool exec) {
L
Liu Jicong 已提交
1279
  ASSERT(0);
1280 1281 1282 1283 1284
  char*              msgStr = pMsg->pCont;
  char*              msgBody = POINTER_SHIFT(msgStr, sizeof(SMsgHead));
  int32_t            msgLen = pMsg->contLen - sizeof(SMsgHead);
  SStreamDispatchReq req;
  SDecoder           decoder;
L
Liu Jicong 已提交
1285
  tDecoderInit(&decoder, (uint8_t*)msgBody, msgLen);
1286
  tDecodeStreamDispatchReq(&decoder, &req);
L
Liu Jicong 已提交
1287 1288 1289 1290
  int32_t taskId = req.taskId;

  SStreamTask* pTask = streamMetaGetTask(pTq->pStreamMeta, taskId);
  if (pTask) {
1291 1292 1293 1294
    SRpcMsg rsp = {
        .info = pMsg->info,
        .code = 0,
    };
L
Liu Jicong 已提交
1295
    streamProcessDispatchReq(pTask, &req, &rsp, exec);
L
Liu Jicong 已提交
1296
    return 0;
1297 1298
  } else {
    return -1;
L
Liu Jicong 已提交
1299
  }
L
Liu Jicong 已提交
1300 1301
}

L
Liu Jicong 已提交
1302
#if 0
L
Liu Jicong 已提交
1303 1304 1305
int32_t tqProcessTaskRecoverReq(STQ* pTq, SRpcMsg* pMsg) {
  SStreamTaskRecoverReq* pReq = pMsg->pCont;
  int32_t                taskId = pReq->taskId;
L
Liu Jicong 已提交
1306 1307 1308
  SStreamTask*           pTask = streamMetaGetTask(pTq->pStreamMeta, taskId);
  if (pTask) {
    streamProcessRecoverReq(pTask, pReq, pMsg);
L
Liu Jicong 已提交
1309
    return 0;
1310 1311
  } else {
    return -1;
L
Liu Jicong 已提交
1312
  }
L
Liu Jicong 已提交
1313 1314
}

L
Liu Jicong 已提交
1315 1316 1317 1318 1319
int32_t tqProcessTaskRecoverRsp(STQ* pTq, SRpcMsg* pMsg) {
  SStreamTaskRecoverRsp* pRsp = pMsg->pCont;
  int32_t                taskId = pRsp->rspTaskId;

  SStreamTask* pTask = streamMetaGetTask(pTq->pStreamMeta, taskId);
L
Liu Jicong 已提交
1320
  if (pTask) {
L
Liu Jicong 已提交
1321
    streamProcessRecoverRsp(pTask, pRsp);
L
Liu Jicong 已提交
1322
    return 0;
1323 1324
  } else {
    return -1;
L
Liu Jicong 已提交
1325
  }
L
Liu Jicong 已提交
1326
}
L
Liu Jicong 已提交
1327
#endif
L
Liu Jicong 已提交
1328

L
Liu Jicong 已提交
1329 1330 1331 1332
int32_t tqProcessTaskDispatchRsp(STQ* pTq, SRpcMsg* pMsg) {
  SStreamDispatchRsp* pRsp = POINTER_SHIFT(pMsg->pCont, sizeof(SMsgHead));
  int32_t             taskId = pRsp->taskId;
  SStreamTask*        pTask = streamMetaGetTask(pTq->pStreamMeta, taskId);
L
Liu Jicong 已提交
1333
  if (pTask) {
L
Liu Jicong 已提交
1334
    streamProcessDispatchRsp(pTask, pRsp);
L
Liu Jicong 已提交
1335
    return 0;
1336 1337
  } else {
    return -1;
L
Liu Jicong 已提交
1338
  }
L
Liu Jicong 已提交
1339
}
L
Liu Jicong 已提交
1340

1341
int32_t tqProcessTaskDropReq(STQ* pTq, int64_t version, char* msg, int32_t msgLen) {
L
Liu Jicong 已提交
1342
  SVDropStreamTaskReq* pReq = (SVDropStreamTaskReq*)msg;
L
Liu Jicong 已提交
1343

L
Liu Jicong 已提交
1344
  return streamMetaRemoveTask(pTq->pStreamMeta, pReq->taskId);
L
Liu Jicong 已提交
1345
}
L
Liu Jicong 已提交
1346 1347 1348 1349 1350 1351 1352 1353 1354

int32_t tqProcessTaskRetrieveReq(STQ* pTq, SRpcMsg* pMsg) {
  char*              msgStr = pMsg->pCont;
  char*              msgBody = POINTER_SHIFT(msgStr, sizeof(SMsgHead));
  int32_t            msgLen = pMsg->contLen - sizeof(SMsgHead);
  SStreamRetrieveReq req;
  SDecoder           decoder;
  tDecoderInit(&decoder, msgBody, msgLen);
  tDecodeStreamRetrieveReq(&decoder, &req);
L
Liu Jicong 已提交
1355
  tDecoderClear(&decoder);
L
Liu Jicong 已提交
1356 1357 1358
  int32_t      taskId = req.dstTaskId;
  SStreamTask* pTask = streamMetaGetTask(pTq->pStreamMeta, taskId);
  if (pTask) {
L
Liu Jicong 已提交
1359 1360 1361 1362
    SRpcMsg rsp = {
        .info = pMsg->info,
        .code = 0,
    };
L
Liu Jicong 已提交
1363
    streamProcessRetrieveReq(pTask, &req, &rsp);
L
Liu Jicong 已提交
1364
    tDeleteStreamRetrieveReq(&req);
L
Liu Jicong 已提交
1365
    return 0;
L
Liu Jicong 已提交
1366 1367
  } else {
    return -1;
L
Liu Jicong 已提交
1368 1369 1370 1371 1372 1373 1374
  }
}

int32_t tqProcessTaskRetrieveRsp(STQ* pTq, SRpcMsg* pMsg) {
  //
  return 0;
}
L
Liu Jicong 已提交
1375 1376 1377 1378 1379 1380 1381 1382 1383 1384 1385 1386 1387

void vnodeEnqueueStreamMsg(SVnode* pVnode, SRpcMsg* pMsg) {
  STQ*    pTq = pVnode->pTq;
  char*   msgStr = pMsg->pCont;
  char*   msgBody = POINTER_SHIFT(msgStr, sizeof(SMsgHead));
  int32_t msgLen = pMsg->contLen - sizeof(SMsgHead);
  int32_t code = 0;

  SStreamDispatchReq req;
  SDecoder           decoder;
  tDecoderInit(&decoder, msgBody, msgLen);
  if (tDecodeStreamDispatchReq(&decoder, &req) < 0) {
    code = TSDB_CODE_MSG_DECODE_ERROR;
L
Liu Jicong 已提交
1388
    tDecoderClear(&decoder);
L
Liu Jicong 已提交
1389 1390
    goto FAIL;
  }
L
Liu Jicong 已提交
1391
  tDecoderClear(&decoder);
L
Liu Jicong 已提交
1392

L
Liu Jicong 已提交
1393
  int32_t taskId = req.taskId;
L
Liu Jicong 已提交
1394

L
Liu Jicong 已提交
1395 1396
  SStreamTask* pTask = streamMetaGetTask(pTq->pStreamMeta, taskId);
  if (pTask) {
L
Liu Jicong 已提交
1397 1398 1399 1400
    SRpcMsg rsp = {
        .info = pMsg->info,
        .code = 0,
    };
L
Liu Jicong 已提交
1401
    streamProcessDispatchReq(pTask, &req, &rsp, false);
L
Liu Jicong 已提交
1402 1403
    rpcFreeCont(pMsg->pCont);
    taosFreeQitem(pMsg);
L
Liu Jicong 已提交
1404 1405
    return;
  }
L
Liu Jicong 已提交
1406

L
Liu Jicong 已提交
1407 1408 1409 1410 1411 1412 1413
FAIL:
  if (pMsg->info.handle == NULL) return;
  SRpcMsg rsp = {
      .code = code,
      .info = pMsg->info,
  };
  tmsgSendRsp(&rsp);
L
Liu Jicong 已提交
1414 1415
  rpcFreeCont(pMsg->pCont);
  taosFreeQitem(pMsg);
L
Liu Jicong 已提交
1416
}