tq.c 45.3 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

dengyihao's avatar
dengyihao 已提交
18 19 20
// 0: not init
// 1: already inited
// 2: wait to be inited or cleaup
21
#define WAL_READ_TASKS_ID       (-1)
22

23 24
static int32_t tqInitialize(STQ* pTq);

L
Liu Jicong 已提交
25
int32_t tqInit() {
L
Liu Jicong 已提交
26 27 28 29 30 31
  int8_t old;
  while (1) {
    old = atomic_val_compare_exchange_8(&tqMgmt.inited, 0, 2);
    if (old != 2) break;
  }

32 33 34 35 36 37
  if (old == 0) {
    tqMgmt.timer = taosTmrInit(10000, 100, 10000, "TQ");
    if (tqMgmt.timer == NULL) {
      atomic_store_8(&tqMgmt.inited, 0);
      return -1;
    }
38 39 40
    if (streamInit() < 0) {
      return -1;
    }
L
Liu Jicong 已提交
41
    atomic_store_8(&tqMgmt.inited, 1);
42
  }
43

L
Liu Jicong 已提交
44 45
  return 0;
}
L
Liu Jicong 已提交
46

47
void tqCleanUp() {
L
Liu Jicong 已提交
48 49 50 51 52 53 54 55
  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 已提交
56
    streamCleanUp();
L
Liu Jicong 已提交
57 58
    atomic_store_8(&tqMgmt.inited, 0);
  }
59
}
L
Liu Jicong 已提交
60

61
static void destroyTqHandle(void* data) {
62 63 64
  STqHandle* pData = (STqHandle*)data;
  qDestroyTask(pData->execHandle.task);
  if (pData->execHandle.subType == TOPIC_SUB_TYPE__COLUMN) {
L
Liu Jicong 已提交
65
    taosMemoryFreeClear(pData->execHandle.execCol.qmsg);
66
  } else if (pData->execHandle.subType == TOPIC_SUB_TYPE__DB) {
67
    tqCloseReader(pData->execHandle.pTqReader);
68 69
    walCloseReader(pData->pWalReader);
    taosHashCleanup(pData->execHandle.execDb.pFilterOutTbUid);
L
Liu Jicong 已提交
70
  } else if (pData->execHandle.subType == TOPIC_SUB_TYPE__TABLE) {
71
    walCloseReader(pData->pWalReader);
72
    tqCloseReader(pData->execHandle.pTqReader);
73 74 75
  }
}

L
Liu Jicong 已提交
76
static void tqPushEntryFree(void* data) {
L
Liu Jicong 已提交
77
  STqPushEntry* p = *(void**)data;
H
Haojun Liao 已提交
78
  if (p->pDataRsp->head.mqMsgType == TMQ_MSG_TYPE__POLL_RSP) {
79
    tDeleteMqDataRsp(p->pDataRsp);
H
Haojun Liao 已提交
80 81 82 83 84
  } else if (p->pDataRsp->head.mqMsgType == TMQ_MSG_TYPE__TAOSX_RSP) {
    tDeleteSTaosxRsp((STaosxRsp*)p->pDataRsp);
  }

  taosMemoryFree(p->pDataRsp);
L
Liu Jicong 已提交
85 86 87
  taosMemoryFree(p);
}

88 89 90 91 92
static 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;
}

L
Liu Jicong 已提交
93
STQ* tqOpen(const char* path, SVnode* pVnode) {
94
  STQ* pTq = taosMemoryCalloc(1, sizeof(STQ));
L
Liu Jicong 已提交
95
  if (pTq == NULL) {
S
Shengliang Guan 已提交
96
    terrno = TSDB_CODE_OUT_OF_MEMORY;
L
Liu Jicong 已提交
97 98
    return NULL;
  }
99

100
  pTq->path = taosStrdup(path);
L
Liu Jicong 已提交
101
  pTq->pVnode = pVnode;
L
Liu Jicong 已提交
102
  pTq->walLogLastVer = pVnode->pWal->vers.lastVer;
103

104
  pTq->pHandle = taosHashInit(64, MurmurHash3_32, true, HASH_ENTRY_LOCK);
105
  taosHashSetFreeFp(pTq->pHandle, destroyTqHandle);
106

107
  taosInitRWLatch(&pTq->lock);
L
Liu Jicong 已提交
108
  pTq->pPushMgr = taosHashInit(64, taosGetDefaultHashFunction(TSDB_DATA_TYPE_BIGINT), true, HASH_NO_LOCK);
L
Liu Jicong 已提交
109
  taosHashSetFreeFp(pTq->pPushMgr, tqPushEntryFree);
L
Liu Jicong 已提交
110

111
  pTq->pCheckInfo = taosHashInit(64, MurmurHash3_32, true, HASH_ENTRY_LOCK);
L
Liu Jicong 已提交
112
  taosHashSetFreeFp(pTq->pCheckInfo, (FDelete)tDeleteSTqCheckInfo);
L
Liu Jicong 已提交
113

114
  tqInitialize(pTq);
115 116 117 118
  return pTq;
}

int32_t tqInitialize(STQ* pTq) {
L
Liu Jicong 已提交
119
  if (tqMetaOpen(pTq) < 0) {
120
    return -1;
121 122
  }

L
Liu Jicong 已提交
123 124
  pTq->pOffsetStore = tqOffsetOpen(pTq);
  if (pTq->pOffsetStore == NULL) {
125
    return -1;
126 127
  }

128
  pTq->pStreamMeta = streamMetaOpen(pTq->path, pTq, (FTaskExpand*)tqExpandTask, pTq->pVnode->config.vgId);
L
Liu Jicong 已提交
129
  if (pTq->pStreamMeta == NULL) {
130
    return -1;
L
Liu Jicong 已提交
131 132
  }

133 134
  // the version is kept in task's meta data
  // todo check if this version is required or not
135 136
  if (streamLoadTasks(pTq->pStreamMeta, walGetCommittedVer(pTq->pVnode->pWal)) < 0) {
    return -1;
L
Liu Jicong 已提交
137 138
  }

139
  return 0;
L
Liu Jicong 已提交
140
}
L
Liu Jicong 已提交
141

L
Liu Jicong 已提交
142
void tqClose(STQ* pTq) {
143 144
  if (pTq == NULL) {
    return;
H
Hongze Cheng 已提交
145
  }
146 147 148 149 150 151 152 153 154

  tqOffsetClose(pTq->pOffsetStore);
  taosHashCleanup(pTq->pHandle);
  taosHashCleanup(pTq->pPushMgr);
  taosHashCleanup(pTq->pCheckInfo);
  taosMemoryFree(pTq->path);
  tqMetaClose(pTq);
  streamMetaClose(pTq->pStreamMeta);
  taosMemoryFree(pTq);
L
Liu Jicong 已提交
155
}
L
Liu Jicong 已提交
156

157
int32_t tqPushDataRsp(STqPushEntry* pPushEntry, int32_t vgId) {
H
Haojun Liao 已提交
158 159
  SMqDataRsp* pRsp = pPushEntry->pDataRsp;
  SMqRspHead* pHeader = &pPushEntry->pDataRsp->head;
160 161 162 163 164

  int64_t sver = 0, ever = 0;
  walReaderValidVersionRange(pPushEntry->pHandle->execHandle.pTqReader->pWalReader, &sver, &ever);

  tqDoSendDataRsp(&pPushEntry->info, pRsp, pHeader->epoch, pHeader->consumerId, pHeader->mqMsgType, sver, ever);
L
Liu Jicong 已提交
165

wmmhello's avatar
wmmhello 已提交
166 167
  char buf1[80] = {0};
  char buf2[80] = {0};
H
Haojun Liao 已提交
168 169 170
  tFormatOffset(buf1, tListLen(buf1), &pRsp->reqOffset);
  tFormatOffset(buf2, tListLen(buf2), &pRsp->rspOffset);
  tqDebug("vgId:%d, from consumer:0x%" PRIx64 " (epoch %d) push rsp, block num: %d, req:%s, rsp:%s",
171
          vgId, pRsp->head.consumerId, pRsp->head.epoch, pRsp->blockNum, buf1, buf2);
L
Liu Jicong 已提交
172 173 174
  return 0;
}

175 176 177 178 179 180
int32_t tqSendDataRsp(STqHandle* pHandle, const SRpcMsg* pMsg, const SMqPollReq* pReq, const SMqDataRsp* pRsp,
                      int32_t type, int32_t vgId) {
  int64_t sver = 0, ever = 0;
  walReaderValidVersionRange(pHandle->execHandle.pTqReader->pWalReader, &sver, &ever);

  tqDoSendDataRsp(&pMsg->info, pRsp, pReq->epoch, pReq->consumerId, type, sver, ever);
181 182 183 184 185 186

  char buf1[80] = {0};
  char buf2[80] = {0};
  tFormatOffset(buf1, 80, &pRsp->reqOffset);
  tFormatOffset(buf2, 80, &pRsp->rspOffset);

X
Xiaoyu Wang 已提交
187
  tqDebug("vgId:%d consumer:0x%" PRIx64 " (epoch %d) send rsp, block num:%d, req:%s, rsp:%s, reqId:0x%" PRIx64,
188
          vgId, pReq->consumerId, pReq->epoch, pRsp->blockNum, buf1, buf2, pReq->reqId);
H
Haojun Liao 已提交
189

190 191 192
  return 0;
}

193
int32_t tqProcessOffsetCommitReq(STQ* pTq, int64_t sversion, char* msg, int32_t msgLen) {
194
  STqOffset offset = {0};
X
Xiaoyu Wang 已提交
195
  int32_t   vgId = TD_VID(pTq->pVnode);
196

X
Xiaoyu Wang 已提交
197 198
  SDecoder decoder;
  tDecoderInit(&decoder, (uint8_t*)msg, msgLen);
199 200 201
  if (tDecodeSTqOffset(&decoder, &offset) < 0) {
    return -1;
  }
202

203 204
  tDecoderClear(&decoder);

wmmhello's avatar
wmmhello 已提交
205
  if (offset.val.type == TMQ_OFFSET__SNAPSHOT_DATA || offset.val.type == TMQ_OFFSET__SNAPSHOT_META) {
L
Liu Jicong 已提交
206
    tqDebug("receive offset commit msg to %s on vgId:%d, offset(type:snapshot) uid:%" PRId64 ", ts:%" PRId64,
H
Haojun Liao 已提交
207
            offset.subKey, vgId, offset.val.uid, offset.val.ts);
L
Liu Jicong 已提交
208
  } else if (offset.val.type == TMQ_OFFSET__LOG) {
X
Xiaoyu Wang 已提交
209 210
    tqDebug("receive offset commit msg to %s on vgId:%d, offset(type:log) version:%" PRId64, offset.subKey, vgId,
            offset.val.version);
211
    if (offset.val.version + 1 == sversion) {
212 213
      offset.val.version += 1;
    }
214
  } else {
215 216
    tqError("invalid commit offset type:%d", offset.val.type);
    return -1;
217
  }
218 219 220

  STqOffset* pSavedOffset = tqOffsetRead(pTq->pOffsetStore, offset.subKey);
  if (pSavedOffset != NULL && tqOffsetLessOrEqual(&offset, pSavedOffset)) {
221 222
    tqDebug("not update the offset, vgId:%d sub:%s since committed:%" PRId64 " less than/equal to existed:%" PRId64,
            vgId, offset.subKey, offset.val.version, pSavedOffset->val.version);
223
    return 0;  // no need to update the offset value
224 225
  }

226
  // save the new offset value
227 228
  if (tqOffsetWrite(pTq->pOffsetStore, &offset) < 0) {
    return -1;
229
  }
230 231

  if (offset.val.type == TMQ_OFFSET__LOG) {
232
    STqHandle* pHandle = taosHashGet(pTq->pHandle, offset.subKey, strlen(offset.subKey));
233 234
    if (pHandle && (walRefVer(pHandle->pRef, offset.val.version) < 0)) {
      return -1;
235 236 237
    }
  }

238 239 240
  return 0;
}

241 242 243 244 245 246 247 248 249 250 251 252 253 254 255 256 257 258 259 260 261 262 263 264 265 266 267 268 269 270 271 272 273 274 275 276 277 278 279 280 281 282 283 284 285 286 287 288 289 290 291
int32_t tqProcessSeekReq(STQ* pTq, int64_t sversion, char* msg, int32_t msgLen) {
  STqOffset offset = {0};
  int32_t   vgId = TD_VID(pTq->pVnode);

  SDecoder decoder;
  tDecoderInit(&decoder, (uint8_t*)msg, msgLen);
  if (tDecodeSTqOffset(&decoder, &offset) < 0) {
    return -1;
  }

  tDecoderClear(&decoder);

  if (offset.val.type != TMQ_OFFSET__LOG) {
    tqError("vgId:%d, subKey:%s invalid seek offset type:%d", vgId, offset.subKey, offset.val.type);
    return -1;
  }

  STqOffset* pSavedOffset = tqOffsetRead(pTq->pOffsetStore, offset.subKey);
  if (pSavedOffset != NULL && pSavedOffset->val.type != TMQ_OFFSET__LOG) {
    tqError("invalid saved offset type, vgId:%d sub:%s", vgId, offset.subKey);
    return 0;  // no need to update the offset value
  }

  if (pSavedOffset->val.version == offset.val.version) {
    tqDebug("vgId:%d subKey:%s no need to seek to %" PRId64 " prev offset:%" PRId64, vgId, offset.subKey, offset.val.version,
            pSavedOffset->val.version);
    return 0;
  }

  STqHandle* pHandle = taosHashGet(pTq->pHandle, offset.subKey, strlen(offset.subKey));

  int64_t sver = 0, ever = 0;
  walReaderValidVersionRange(pHandle->execHandle.pTqReader->pWalReader, &sver, &ever);
  if (offset.val.version < sver) {
    offset.val.version = sver;
  } else if (offset.val.version > ever) {
    offset.val.version = ever;
  }

  // save the new offset value
  tqDebug("vgId:%d sub:%s seek to %" PRId64 " prev offset:%" PRId64, vgId, offset.subKey, offset.val.version,
          pSavedOffset->val.version);

  if (tqOffsetWrite(pTq->pOffsetStore, &offset) < 0) {
    tqError("failed to save offset, vgId:%d sub:%s seek to %" PRId64, vgId, offset.subKey, offset.val.version);
    return -1;
  }

  return 0;
}

L
Liu Jicong 已提交
292
int32_t tqCheckColModifiable(STQ* pTq, int64_t tbUid, int32_t colId) {
L
Liu Jicong 已提交
293
  void* pIter = NULL;
294

L
Liu Jicong 已提交
295
  while (1) {
296
    pIter = taosHashIterate(pTq->pCheckInfo, pIter);
297 298 299 300
    if (pIter == NULL) {
      break;
    }

301
    STqCheckInfo* pCheck = (STqCheckInfo*)pIter;
302

L
Liu Jicong 已提交
303 304
    if (pCheck->ntbUid == tbUid) {
      int32_t sz = taosArrayGetSize(pCheck->colIdList);
L
Liu Jicong 已提交
305
      for (int32_t i = 0; i < sz; i++) {
L
Liu Jicong 已提交
306 307
        int16_t forbidColId = *(int16_t*)taosArrayGet(pCheck->colIdList, i);
        if (forbidColId == colId) {
308
          taosHashCancelIterate(pTq->pCheckInfo, pIter);
L
Liu Jicong 已提交
309 310 311 312 313
          return -1;
        }
      }
    }
  }
314

L
Liu Jicong 已提交
315 316 317
  return 0;
}

318
int32_t tqProcessPollReq(STQ* pTq, SRpcMsg* pMsg) {
X
Xiaoyu Wang 已提交
319
  SMqPollReq req = {0};
320 321 322 323 324 325 326 327 328 329 330 331 332 333 334 335 336 337 338
  if (tDeserializeSMqPollReq(pMsg->pCont, pMsg->contLen, &req) < 0) {
    tqError("tDeserializeSMqPollReq %d failed", pMsg->contLen);
    terrno = TSDB_CODE_INVALID_MSG;
    return -1;
  }

  int64_t      consumerId = req.consumerId;
  int32_t      reqEpoch = req.epoch;
  STqOffsetVal reqOffset = req.reqOffset;
  int32_t      vgId = TD_VID(pTq->pVnode);

  // 1. find handle
  STqHandle* pHandle = taosHashGet(pTq->pHandle, req.subKey, strlen(req.subKey));
  if (pHandle == NULL) {
    tqError("tmq poll: consumer:0x%" PRIx64 " vgId:%d subkey %s not found", consumerId, vgId, req.subKey);
    terrno = TSDB_CODE_INVALID_MSG;
    return -1;
  }

339
  // 2. check re-balance status
340
  taosRLockLatch(&pTq->lock);
341 342 343 344
  if (pHandle->consumerId != consumerId) {
    tqDebug("ERROR tmq poll: consumer:0x%" PRIx64 " vgId:%d, subkey %s, mismatch for saved handle consumer:0x%" PRIx64,
            consumerId, TD_VID(pTq->pVnode), req.subKey, pHandle->consumerId);
    terrno = TSDB_CODE_TMQ_CONSUMER_MISMATCH;
345
    taosRUnLockLatch(&pTq->lock);
346 347
    return -1;
  }
348
  taosRUnLockLatch(&pTq->lock);
349

350
  // 3. update the epoch value
351
  taosWLockLatch(&pTq->lock);
H
Haojun Liao 已提交
352 353
  int32_t savedEpoch = pHandle->epoch;
  if (savedEpoch < reqEpoch) {
X
Xiaoyu Wang 已提交
354 355
    tqDebug("tmq poll: consumer:0x%" PRIx64 " epoch update from %d to %d by poll req", consumerId, savedEpoch,
            reqEpoch);
356
    pHandle->epoch = reqEpoch;
H
Haojun Liao 已提交
357
  }
358
  taosWUnLockLatch(&pTq->lock);
359 360 361

  char buf[80];
  tFormatOffset(buf, 80, &reqOffset);
H
Haojun Liao 已提交
362 363
  tqDebug("tmq poll: consumer:0x%" PRIx64 " (epoch %d), subkey %s, recv poll req vgId:%d, req:%s, reqId:0x%" PRIx64,
          consumerId, req.epoch, pHandle->subKey, vgId, buf, req.reqId);
364

365
  return tqExtractDataForMq(pTq, pHandle, &req, pMsg);
366 367
}

368 369 370 371 372 373 374 375 376 377 378 379 380 381 382 383 384 385 386 387 388 389 390 391 392 393 394 395 396 397 398 399 400 401 402 403 404 405 406 407 408 409 410 411 412 413 414 415 416 417 418 419 420 421 422 423 424 425 426 427 428 429 430 431 432 433 434 435 436 437 438 439 440 441 442
int32_t tqProcessVgWalInfoReq(STQ* pTq, SRpcMsg* pMsg) {
  SMqPollReq req = {0};
  if (tDeserializeSMqPollReq(pMsg->pCont, pMsg->contLen, &req) < 0) {
    tqError("tDeserializeSMqPollReq %d failed", pMsg->contLen);
    terrno = TSDB_CODE_INVALID_MSG;
    return -1;
  }

  int64_t      consumerId = req.consumerId;
  STqOffsetVal reqOffset = req.reqOffset;
  int32_t      vgId = TD_VID(pTq->pVnode);

  // 1. find handle
  STqHandle* pHandle = taosHashGet(pTq->pHandle, req.subKey, strlen(req.subKey));
  if (pHandle == NULL) {
    tqError("consumer:0x%" PRIx64 " vgId:%d subkey:%s not found", consumerId, vgId, req.subKey);
    terrno = TSDB_CODE_INVALID_MSG;
    return -1;
  }

  // 2. check re-balance status
  taosRLockLatch(&pTq->lock);
  if (pHandle->consumerId != consumerId) {
    tqDebug("ERROR consumer:0x%" PRIx64 " vgId:%d, subkey %s, mismatch for saved handle consumer:0x%" PRIx64,
            consumerId, vgId, req.subKey, pHandle->consumerId);
    terrno = TSDB_CODE_TMQ_CONSUMER_MISMATCH;
    taosRUnLockLatch(&pTq->lock);
    return -1;
  }
  taosRUnLockLatch(&pTq->lock);

  int64_t sver = 0, ever = 0;
  walReaderValidVersionRange(pHandle->execHandle.pTqReader->pWalReader, &sver, &ever);

  SMqDataRsp dataRsp = {0};
  tqInitDataRsp(&dataRsp, &req);

  STqOffset* pOffset = tqOffsetRead(pTq->pOffsetStore, req.subKey);
  if (pOffset != NULL) {
    if (pOffset->val.type != TMQ_OFFSET__LOG) {
      tqError("consumer:0x%" PRIx64 " vgId:%d subkey:%s use snapshot, no valid wal info", consumerId, vgId, req.subKey);
      terrno = TSDB_CODE_INVALID_PARA;
      tDeleteMqDataRsp(&dataRsp);
      return -1;
    }

    dataRsp.rspOffset.type = TMQ_OFFSET__LOG;
    dataRsp.rspOffset.version = pOffset->val.version;
  } else {
    if (req.useSnapshot == true) {
      tqError("consumer:0x%" PRIx64 " vgId:%d subkey:%s snapshot not support wal info", consumerId, vgId, req.subKey);
      terrno = TSDB_CODE_INVALID_PARA;
      tDeleteMqDataRsp(&dataRsp);
      return -1;
    }

    dataRsp.rspOffset.type = TMQ_OFFSET__LOG;

    if (reqOffset.type == TMQ_OFFSET__RESET_EARLIEAST) {
      dataRsp.rspOffset.version = sver;
    } else if (reqOffset.type == TMQ_OFFSET__RESET_LATEST) {
      dataRsp.rspOffset.version = ever;
    } else {
      tqError("consumer:0x%" PRIx64 " vgId:%d subkey:%s invalid offset type:%d", consumerId, vgId, req.subKey,
              reqOffset.type);
      terrno = TSDB_CODE_INVALID_PARA;
      tDeleteMqDataRsp(&dataRsp);
      return -1;
    }
  }

  tqDoSendDataRsp(&pMsg->info, &dataRsp, req.epoch, req.consumerId, TMQ_MSG_TYPE__WALINFO_RSP, sver, ever);
  return 0;
}

443
int32_t tqProcessDeleteSubReq(STQ* pTq, int64_t sversion, char* msg, int32_t msgLen) {
L
Liu Jicong 已提交
444
  SMqVDeleteReq* pReq = (SMqVDeleteReq*)msg;
L
Liu Jicong 已提交
445

L
Liu Jicong 已提交
446
  tqDebug("vgId:%d, tq process delete sub req %s", pTq->pVnode->config.vgId, pReq->subKey);
L
Liu Jicong 已提交
447

448
  taosWLockLatch(&pTq->lock);
L
Liu Jicong 已提交
449 450 451 452
  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);
  }
453
  taosWUnLockLatch(&pTq->lock);
L
Liu Jicong 已提交
454

L
Liu Jicong 已提交
455 456
  STqHandle* pHandle = taosHashGet(pTq->pHandle, pReq->subKey, strlen(pReq->subKey));
  if (pHandle) {
X
Xiaoyu Wang 已提交
457
    // walCloseRef(pHandle->pWalReader->pWal, pHandle->pRef->refId);
L
Liu Jicong 已提交
458 459 460 461 462 463 464
    if (pHandle->pRef) {
      walCloseRef(pTq->pVnode->pWal, pHandle->pRef->refId);
    }
    code = taosHashRemove(pTq->pHandle, pReq->subKey, strlen(pReq->subKey));
    if (code != 0) {
      tqError("cannot process tq delete req %s, since no such handle", pReq->subKey);
    }
L
Liu Jicong 已提交
465
  }
466

L
Liu Jicong 已提交
467 468
  code = tqOffsetDelete(pTq->pOffsetStore, pReq->subKey);
  if (code != 0) {
469
    tqError("cannot process tq delete req %s, since no such offset in cache", pReq->subKey);
L
Liu Jicong 已提交
470
  }
L
Liu Jicong 已提交
471

L
Liu Jicong 已提交
472
  if (tqMetaDeleteHandle(pTq, pReq->subKey) < 0) {
L
Liu Jicong 已提交
473
    tqError("cannot process tq delete req %s, since no such offset in tdb", pReq->subKey);
474
  }
L
Liu Jicong 已提交
475
  return 0;
L
Liu Jicong 已提交
476 477
}

478
int32_t tqProcessAddCheckInfoReq(STQ* pTq, int64_t sversion, char* msg, int32_t msgLen) {
479 480
  STqCheckInfo info = {0};
  SDecoder     decoder;
X
Xiaoyu Wang 已提交
481
  tDecoderInit(&decoder, (uint8_t*)msg, msgLen);
482
  if (tDecodeSTqCheckInfo(&decoder, &info) < 0) {
L
Liu Jicong 已提交
483 484 485 486
    terrno = TSDB_CODE_OUT_OF_MEMORY;
    return -1;
  }
  tDecoderClear(&decoder);
487 488 489 490 491
  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 已提交
492 493 494 495 496 497
    terrno = TSDB_CODE_OUT_OF_MEMORY;
    return -1;
  }
  return 0;
}

498
int32_t tqProcessDelCheckInfoReq(STQ* pTq, int64_t sversion, char* msg, int32_t msgLen) {
499 500 501 502 503 504 505 506 507 508 509
  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;
}

510
int32_t tqProcessSubscribeReq(STQ* pTq, int64_t sversion, char* msg, int32_t msgLen) {
L
Liu Jicong 已提交
511
  SMqRebVgReq req = {0};
L
Liu Jicong 已提交
512
  tDecodeSMqRebVgReq(msg, &req);
L
Liu Jicong 已提交
513

514 515 516
  SVnode* pVnode = pTq->pVnode;
  int32_t vgId = TD_VID(pVnode);

517
  tqDebug("vgId:%d, tq process sub req:%s, Id:0x%" PRIx64 " -> Id:0x%" PRIx64, pVnode->config.vgId, req.subKey,
518
          req.oldConsumerId, req.newConsumerId);
L
Liu Jicong 已提交
519

520
  STqHandle* pHandle = taosHashGet(pTq->pHandle, req.subKey, strlen(req.subKey));
L
Liu Jicong 已提交
521
  if (pHandle == NULL) {
L
Liu Jicong 已提交
522
    if (req.oldConsumerId != -1) {
523
      tqError("vgId:%d, build new consumer handle %s for consumer:0x%" PRIx64 ", but old consumerId:0x%" PRIx64,
524
              req.vgId, req.subKey, req.newConsumerId, req.oldConsumerId);
L
Liu Jicong 已提交
525
    }
526

L
Liu Jicong 已提交
527
    if (req.newConsumerId == -1) {
528
      tqError("vgId:%d, tq invalid re-balance request, new consumerId %" PRId64 "", req.vgId, req.newConsumerId);
L
Liu Jicong 已提交
529
      taosMemoryFree(req.qmsg);
L
Liu Jicong 已提交
530 531
      return 0;
    }
532

L
Liu Jicong 已提交
533 534
    STqHandle tqHandle = {0};
    pHandle = &tqHandle;
L
Liu Jicong 已提交
535

H
Haojun Liao 已提交
536
    uint64_t oldConsumerId = pHandle->consumerId;
L
Liu Jicong 已提交
537 538 539
    memcpy(pHandle->subKey, req.subKey, TSDB_SUBSCRIBE_KEY_LEN);
    pHandle->consumerId = req.newConsumerId;
    pHandle->epoch = -1;
L
Liu Jicong 已提交
540

L
Liu Jicong 已提交
541
    pHandle->execHandle.subType = req.subType;
L
Liu Jicong 已提交
542
    pHandle->fetchMeta = req.withMeta;
wmmhello's avatar
wmmhello 已提交
543

544
    // TODO version should be assigned and refed during preprocess
545
    SWalRef* pRef = walRefCommittedVer(pVnode->pWal);
546
    if (pRef == NULL) {
H
Haojun Liao 已提交
547
      taosMemoryFree(req.qmsg);
L
Liu Jicong 已提交
548
      return -1;
549
    }
H
Haojun Liao 已提交
550

551 552
    int64_t ver = pRef->refVer;
    pHandle->pRef = pRef;
L
Liu Jicong 已提交
553

554
    SReadHandle handle = {
555
        .meta = pVnode->pMeta, .vnode = pVnode, .initTableReader = true, .initTqReader = true, .version = ver};
wmmhello's avatar
wmmhello 已提交
556
    pHandle->snapshotVer = ver;
557

L
Liu Jicong 已提交
558
    if (pHandle->execHandle.subType == TOPIC_SUB_TYPE__COLUMN) {
L
Liu Jicong 已提交
559
      pHandle->execHandle.execCol.qmsg = req.qmsg;
L
Liu Jicong 已提交
560
      req.qmsg = NULL;
561

X
Xiaoyu Wang 已提交
562 563
      pHandle->execHandle.task = qCreateQueueExecTaskInfo(pHandle->execHandle.execCol.qmsg, &handle, vgId,
                                                          &pHandle->execHandle.numOfCols, req.newConsumerId);
L
Liu Jicong 已提交
564
      void* scanner = NULL;
565
      qExtractStreamScanner(pHandle->execHandle.task, &scanner);
566
      pHandle->execHandle.pTqReader = qExtractReaderFromStreamScanner(scanner);
L
Liu Jicong 已提交
567
    } else if (pHandle->execHandle.subType == TOPIC_SUB_TYPE__DB) {
568
      pHandle->pWalReader = walOpenReader(pVnode->pWal, NULL);
569
      pHandle->execHandle.pTqReader = tqOpenReader(pVnode);
570

L
Liu Jicong 已提交
571
      pHandle->execHandle.execDb.pFilterOutTbUid =
L
Liu Jicong 已提交
572
          taosHashInit(64, taosGetDefaultHashFunction(TSDB_DATA_TYPE_BIGINT), false, HASH_NO_LOCK);
573 574
      buildSnapContext(handle.meta, handle.version, 0, pHandle->execHandle.subType, pHandle->fetchMeta,
                       (SSnapContext**)(&handle.sContext));
575

576
      pHandle->execHandle.task = qCreateQueueExecTaskInfo(NULL, &handle, vgId, NULL, req.newConsumerId);
L
Liu Jicong 已提交
577
    } else if (pHandle->execHandle.subType == TOPIC_SUB_TYPE__TABLE) {
578
      pHandle->pWalReader = walOpenReader(pVnode->pWal, NULL);
wmmhello's avatar
wmmhello 已提交
579 580
      pHandle->execHandle.execTb.suid = req.suid;

L
Liu Jicong 已提交
581
      SArray* tbUidList = taosArrayInit(0, sizeof(int64_t));
582 583
      vnodeGetCtbIdList(pVnode, req.suid, tbUidList);
      tqDebug("vgId:%d, tq try to get all ctb, suid:%" PRId64, pVnode->config.vgId, req.suid);
L
Liu Jicong 已提交
584 585
      for (int32_t i = 0; i < taosArrayGetSize(tbUidList); i++) {
        int64_t tbUid = *(int64_t*)taosArrayGet(tbUidList, i);
586
        tqDebug("vgId:%d, idx %d, uid:%" PRId64, vgId, i, tbUid);
L
Liu Jicong 已提交
587
      }
588 589
      pHandle->execHandle.pTqReader = tqOpenReader(pVnode);
      tqReaderSetTbUidList(pHandle->execHandle.pTqReader, tbUidList);
L
Liu Jicong 已提交
590
      taosArrayDestroy(tbUidList);
wmmhello's avatar
wmmhello 已提交
591

L
Liu Jicong 已提交
592 593
      buildSnapContext(handle.meta, handle.version, req.suid, pHandle->execHandle.subType, pHandle->fetchMeta,
                       (SSnapContext**)(&handle.sContext));
594
      pHandle->execHandle.task = qCreateQueueExecTaskInfo(NULL, &handle, vgId, NULL, req.newConsumerId);
L
Liu Jicong 已提交
595
    }
H
Haojun Liao 已提交
596

597
    taosHashPut(pTq->pHandle, req.subKey, strlen(req.subKey), pHandle, sizeof(STqHandle));
598 599
    tqDebug("try to persist handle %s consumer:0x%" PRIx64 " , old consumer:0x%" PRIx64, req.subKey,
            pHandle->consumerId, oldConsumerId);
L
Liu Jicong 已提交
600
    if (tqMetaSaveHandle(pTq, req.subKey, pHandle) < 0) {
H
Haojun Liao 已提交
601
      taosMemoryFree(req.qmsg);
L
Liu Jicong 已提交
602
      return -1;
L
Liu Jicong 已提交
603
    }
L
Liu Jicong 已提交
604
  } else {
605 606 607 608
    if (pHandle->consumerId == req.newConsumerId) {  // do nothing
      tqInfo("vgId:%d consumer:0x%" PRIx64 " remains, no switch occurs", req.vgId, req.newConsumerId);
      atomic_store_32(&pHandle->epoch, -1);
      atomic_add_fetch_32(&pHandle->epoch, 1);
H
Haojun Liao 已提交
609
      taosMemoryFree(req.qmsg);
610
      return tqMetaSaveHandle(pTq, req.subKey, pHandle);
611 612 613
    } else {
      tqInfo("vgId:%d switch consumer from Id:0x%" PRIx64 " to Id:0x%" PRIx64, req.vgId, pHandle->consumerId,
             req.newConsumerId);
614

615 616 617
      // kill executing task
      qTaskInfo_t pTaskInfo = pHandle->execHandle.task;
      if (pTaskInfo != NULL) {
618
        qKillTask(pTaskInfo, TSDB_CODE_SUCCESS);
619 620
      }

621 622
      taosWLockLatch(&pTq->lock);
      atomic_store_32(&pHandle->epoch, -1);
623

624
      // remove if it has been register in the push manager, and return one empty block to consumer
625
      tqUnregisterPushHandle(pTq, req.subKey, (int32_t)strlen(req.subKey), pHandle->consumerId, true);
626

627 628
      atomic_store_64(&pHandle->consumerId, req.newConsumerId);
      atomic_add_fetch_32(&pHandle->epoch, 1);
629

630
      if (pHandle->execHandle.subType == TOPIC_SUB_TYPE__COLUMN) {
631 632 633
        qStreamCloseTsdbReader(pTaskInfo);
      }

634 635 636 637 638
      taosWUnLockLatch(&pTq->lock);
      if (tqMetaSaveHandle(pTq, req.subKey, pHandle) < 0) {
        taosMemoryFree(req.qmsg);
        return -1;
      }
L
Liu Jicong 已提交
639
    }
L
Liu Jicong 已提交
640
  }
L
Liu Jicong 已提交
641

H
Haojun Liao 已提交
642
  taosMemoryFree(req.qmsg);
L
Liu Jicong 已提交
643
  return 0;
L
Liu Jicong 已提交
644
}
645

646
int32_t tqExpandTask(STQ* pTq, SStreamTask* pTask, int64_t ver) {
647
  int32_t vgId = TD_VID(pTq->pVnode);
648
  pTask->id.idStr = createStreamTaskIdStr(pTask->id.streamId, pTask->id.taskId);
L
Liu Jicong 已提交
649
  pTask->refCnt = 1;
650
  pTask->status.schedStatus = TASK_SCHED_STATUS__INACTIVE;
L
Liu Jicong 已提交
651 652
  pTask->inputQueue = streamQueueOpen();
  pTask->outputQueue = streamQueueOpen();
L
Liu Jicong 已提交
653 654

  if (pTask->inputQueue == NULL || pTask->outputQueue == NULL) {
L
Liu Jicong 已提交
655
    return -1;
L
Liu Jicong 已提交
656 657
  }

L
Liu Jicong 已提交
658 659
  pTask->inputStatus = TASK_INPUT_STATUS__NORMAL;
  pTask->outputStatus = TASK_OUTPUT_STATUS__NORMAL;
660
  pTask->pMsgCb = &pTq->pVnode->msgCb;
661
  pTask->pMeta = pTq->pStreamMeta;
662
  pTask->chkInfo.version = ver;
663

664
  // expand executor
665
  if (pTask->fillHistory) {
666
    pTask->status.taskStatus = TASK_STATUS__WAIT_DOWNSTREAM;
667
  } else {
668
    pTask->status.taskStatus = TASK_STATUS__RESTORE;
669 670
  }

671
  if (pTask->taskLevel == TASK_LEVEL__SOURCE) {
672
    pTask->pState = streamStateOpen(pTq->pStreamMeta->path, pTask, false, -1, -1);
673 674 675 676
    if (pTask->pState == NULL) {
      return -1;
    }

677
    SReadHandle handle = {
678
        .meta = pTq->pVnode->pMeta, .vnode = pTq->pVnode, .initTqReader = 1, .pStateBackend = pTask->pState};
679

680 681
    pTask->exec.pExecutor = qCreateStreamExecTaskInfo(pTask->exec.qmsg, &handle, vgId);
    if (pTask->exec.pExecutor == NULL) {
L
Liu Jicong 已提交
682 683
      return -1;
    }
684

685
  } else if (pTask->taskLevel == TASK_LEVEL__AGG) {
686
    pTask->pState = streamStateOpen(pTq->pStreamMeta->path, pTask, false, -1, -1);
687 688 689
    if (pTask->pState == NULL) {
      return -1;
    }
690

691 692 693 694 695
    int32_t numOfVgroups = (int32_t)taosArrayGetSize(pTask->childEpInfo);
    SReadHandle mgHandle = { .vnode = NULL, .numOfVgroups = numOfVgroups, .pStateBackend = pTask->pState};

    pTask->exec.pExecutor = qCreateStreamExecTaskInfo(pTask->exec.qmsg, &mgHandle, vgId);
    if (pTask->exec.pExecutor == NULL) {
L
Liu Jicong 已提交
696 697
      return -1;
    }
L
Liu Jicong 已提交
698
  }
L
Liu Jicong 已提交
699 700

  // sink
L
Liu Jicong 已提交
701
  /*pTask->ahandle = pTq->pVnode;*/
702
  if (pTask->outputType == TASK_OUTPUT__SMA) {
L
Liu Jicong 已提交
703
    pTask->smaSink.vnode = pTq->pVnode;
L
Liu Jicong 已提交
704
    pTask->smaSink.smaSink = smaHandleRes;
705
  } else if (pTask->outputType == TASK_OUTPUT__TABLE) {
L
Liu Jicong 已提交
706
    pTask->tbSink.vnode = pTq->pVnode;
L
Liu Jicong 已提交
707
    pTask->tbSink.tbSinkFunc = tqSinkToTablePipeline2;
L
Liu Jicong 已提交
708

X
Xiaoyu Wang 已提交
709
    int32_t   ver1 = 1;
5
54liuyao 已提交
710
    SMetaInfo info = {0};
dengyihao's avatar
dengyihao 已提交
711
    int32_t   code = metaGetInfo(pTq->pVnode->pMeta, pTask->tbSink.stbUid, &info, NULL);
5
54liuyao 已提交
712
    if (code == TSDB_CODE_SUCCESS) {
H
Haojun Liao 已提交
713
      ver1 = info.skmVer;
5
54liuyao 已提交
714
    }
L
Liu Jicong 已提交
715

716 717
    SSchemaWrapper* pschemaWrapper = pTask->tbSink.pSchemaWrapper;
    pTask->tbSink.pTSchema = tBuildTSchema(pschemaWrapper->pSchema, pschemaWrapper->nCols, ver1);
wmmhello's avatar
wmmhello 已提交
718
    if(pTask->tbSink.pTSchema == NULL) {
wmmhello's avatar
wmmhello 已提交
719
      return -1;
wmmhello's avatar
wmmhello 已提交
720
    }
L
Liu Jicong 已提交
721
  }
722

723
  if (pTask->taskLevel == TASK_LEVEL__SOURCE) {
724
    pTask->exec.pWalReader = walOpenReader(pTq->pVnode->pWal, NULL);
725 726
  }

727
  streamSetupTrigger(pTask);
728
  tqInfo("vgId:%d expand stream task, s-task:%s, checkpoint ver:%" PRId64 " child id:%d, level:%d", vgId, pTask->id.idStr,
729
         pTask->chkInfo.version, pTask->selfChildId, pTask->taskLevel);
730 731 732

  // next valid version will add one
  pTask->chkInfo.version += 1;
L
Liu Jicong 已提交
733
  return 0;
L
Liu Jicong 已提交
734
}
L
Liu Jicong 已提交
735

736 737 738 739 740 741
int32_t tqProcessStreamTaskCheckReq(STQ* pTq, SRpcMsg* pMsg) {
  char*               msgStr = pMsg->pCont;
  char*               msgBody = POINTER_SHIFT(msgStr, sizeof(SMsgHead));
  int32_t             msgLen = pMsg->contLen - sizeof(SMsgHead);
  SStreamTaskCheckReq req;
  SDecoder            decoder;
X
Xiaoyu Wang 已提交
742
  tDecoderInit(&decoder, (uint8_t*)msgBody, msgLen);
743 744 745 746 747 748 749 750 751 752 753 754
  tDecodeSStreamTaskCheckReq(&decoder, &req);
  tDecoderClear(&decoder);
  int32_t             taskId = req.downstreamTaskId;
  SStreamTaskCheckRsp rsp = {
      .reqId = req.reqId,
      .streamId = req.streamId,
      .childId = req.childId,
      .downstreamNodeId = req.downstreamNodeId,
      .downstreamTaskId = req.downstreamTaskId,
      .upstreamNodeId = req.upstreamNodeId,
      .upstreamTaskId = req.upstreamTaskId,
  };
755

L
Liu Jicong 已提交
756
  SStreamTask* pTask = streamMetaAcquireTask(pTq->pStreamMeta, taskId);
757 758 759 760 761 762 763 764
  if (pTask) {
    rsp.status = (atomic_load_8(&pTask->status.taskStatus) == TASK_STATUS__NORMAL) ? 1 : 0;
    streamMetaReleaseTask(pTq->pStreamMeta, pTask);

    tqDebug("tq recv task check req(reqId:0x%" PRIx64
            ") %d at node %d task status:%d, check req from task %d at node %d, rsp status %d",
            rsp.reqId, rsp.downstreamTaskId, rsp.downstreamNodeId, pTask->status.taskStatus, rsp.upstreamTaskId,
            rsp.upstreamNodeId, rsp.status);
765 766
  } else {
    rsp.status = 0;
767 768 769 770
    tqDebug("tq recv task check(taskId:%d not built yet) req(reqId:0x%" PRIx64
            ") %d at node %d, check req from task %d at node %d, rsp status %d",
            taskId, rsp.reqId, rsp.downstreamTaskId, rsp.downstreamNodeId, rsp.upstreamTaskId, rsp.upstreamNodeId,
            rsp.status);
771 772 773 774 775 776 777
  }

  SEncoder encoder;
  int32_t  code;
  int32_t  len;
  tEncodeSize(tEncodeSStreamTaskCheckRsp, &rsp, len, code);
  if (code < 0) {
L
Liu Jicong 已提交
778
    tqError("unable to encode rsp %d", __LINE__);
L
Liu Jicong 已提交
779
    return -1;
780
  }
L
Liu Jicong 已提交
781

782 783 784 785 786 787 788 789
  void* buf = rpcMallocCont(sizeof(SMsgHead) + len);
  ((SMsgHead*)buf)->vgId = htonl(req.upstreamNodeId);

  void* abuf = POINTER_SHIFT(buf, sizeof(SMsgHead));
  tEncoderInit(&encoder, (uint8_t*)abuf, len);
  tEncodeSStreamTaskCheckRsp(&encoder, &rsp);
  tEncoderClear(&encoder);

790
  SRpcMsg rspMsg = { .code = 0, .pCont = buf, .contLen = sizeof(SMsgHead) + len, .info = pMsg->info };
791 792 793 794
  tmsgSendRsp(&rspMsg);
  return 0;
}

795
int32_t tqProcessStreamTaskCheckRsp(STQ* pTq, int64_t sversion, char* msg, int32_t msgLen) {
796 797 798 799 800 801 802 803 804 805 806
  int32_t             code;
  SStreamTaskCheckRsp rsp;

  SDecoder decoder;
  tDecoderInit(&decoder, (uint8_t*)msg, msgLen);
  code = tDecodeSStreamTaskCheckRsp(&decoder, &rsp);
  if (code < 0) {
    tDecoderClear(&decoder);
    return -1;
  }

807
  tDecoderClear(&decoder);
808
  tqDebug("tq recv task check rsp(reqId:0x%" PRIx64 ") %d at node %d check req from task %d at node %d, status %d",
809 810
          rsp.reqId, rsp.downstreamTaskId, rsp.downstreamNodeId, rsp.upstreamTaskId, rsp.upstreamNodeId, rsp.status);

L
Liu Jicong 已提交
811
  SStreamTask* pTask = streamMetaAcquireTask(pTq->pStreamMeta, rsp.upstreamTaskId);
812 813 814 815
  if (pTask == NULL) {
    return -1;
  }

816
  code = streamProcessTaskCheckRsp(pTask, &rsp, sversion);
L
Liu Jicong 已提交
817 818
  streamMetaReleaseTask(pTq->pStreamMeta, pTask);
  return code;
819 820
}

821
int32_t tqProcessTaskDeployReq(STQ* pTq, int64_t sversion, char* msg, int32_t msgLen) {
822 823 824 825 826
  int32_t code;
#if 0
  code = streamMetaAddSerializedTask(pTq->pStreamMeta, version, msg, msgLen);
  if (code < 0) return code;
#endif
5
54liuyao 已提交
827 828 829
  if (tsDisableStream) {
    return 0;
  }
830 831 832 833 834 835

  // 1.deserialize msg and build task
  SStreamTask* pTask = taosMemoryCalloc(1, sizeof(SStreamTask));
  if (pTask == NULL) {
    return -1;
  }
836

837 838
  SDecoder decoder;
  tDecoderInit(&decoder, (uint8_t*)msg, msgLen);
839
  code = tDecodeStreamTask(&decoder, pTask);
840 841 842 843 844
  if (code < 0) {
    tDecoderClear(&decoder);
    taosMemoryFree(pTask);
    return -1;
  }
845

846 847
  tDecoderClear(&decoder);

848
  // 2.save task, use the newest commit version as the initial start version of stream task.
849
  code = streamMetaAddDeployedTask(pTq->pStreamMeta, sversion, pTask);
850
  if (code < 0) {
851 852
    tqError("vgId:%d failed to add s-task:%s, total:%d", TD_VID(pTq->pVnode), pTask->id.idStr,
            streamMetaGetNumOfTasks(pTq->pStreamMeta));
853 854 855 856 857
    return -1;
  }

  // 3.go through recover steps to fill history
  if (pTask->fillHistory) {
858
    streamTaskCheckDownstream(pTask, sversion);
859 860
  }

861 862
  tqDebug("vgId:%d s-task:%s is deployed and add meta from mnd, status:%d, total:%d", TD_VID(pTq->pVnode),
          pTask->id.idStr, pTask->status.taskStatus, streamMetaGetNumOfTasks(pTq->pStreamMeta));
863 864 865
  return 0;
}

L
Liu Jicong 已提交
866 867 868 869 870
int32_t tqProcessTaskRecover1Req(STQ* pTq, SRpcMsg* pMsg) {
  int32_t code;
  char*   msg = pMsg->pCont;
  int32_t msgLen = pMsg->contLen;

871
  SStreamRecoverStep1Req* pReq = (SStreamRecoverStep1Req*)msg;
L
Liu Jicong 已提交
872
  SStreamTask*            pTask = streamMetaAcquireTask(pTq->pStreamMeta, pReq->taskId);
873 874 875 876 877
  if (pTask == NULL) {
    return -1;
  }

  // check param
878
  int64_t fillVer1 = pTask->chkInfo.version;
879
  if (fillVer1 <= 0) {
L
Liu Jicong 已提交
880
    streamMetaReleaseTask(pTq->pStreamMeta, pTask);
881 882 883 884 885 886
    return -1;
  }

  // do recovery step 1
  streamSourceRecoverScanStep1(pTask);

887
  if (atomic_load_8(&pTask->status.taskStatus) == TASK_STATUS__DROPPING) {
L
Liu Jicong 已提交
888 889 890 891
    streamMetaReleaseTask(pTq->pStreamMeta, pTask);
    return 0;
  }

892 893 894 895
  // build msg to launch next step
  SStreamRecoverStep2Req req;
  code = streamBuildSourceRecover2Req(pTask, &req);
  if (code < 0) {
L
Liu Jicong 已提交
896
    streamMetaReleaseTask(pTq->pStreamMeta, pTask);
897 898 899
    return -1;
  }

L
Liu Jicong 已提交
900
  streamMetaReleaseTask(pTq->pStreamMeta, pTask);
L
Liu Jicong 已提交
901

902
  if (atomic_load_8(&pTask->status.taskStatus) == TASK_STATUS__DROPPING) {
L
Liu Jicong 已提交
903 904 905
    return 0;
  }

906
  // serialize msg
L
Liu Jicong 已提交
907 908 909 910 911 912 913 914
  int32_t len = sizeof(SStreamRecoverStep1Req);

  void* serializedReq = rpcMallocCont(len);
  if (serializedReq == NULL) {
    return -1;
  }

  memcpy(serializedReq, &req, len);
915 916 917 918 919

  // dispatch msg
  SRpcMsg rpcMsg = {
      .code = 0,
      .contLen = len,
L
Liu Jicong 已提交
920
      .msgType = TDMT_VND_STREAM_RECOVER_BLOCKING_STAGE,
L
Liu Jicong 已提交
921
      .pCont = serializedReq,
922 923 924 925 926 927 928
  };

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

  return 0;
}

929
int32_t tqProcessTaskRecover2Req(STQ* pTq, int64_t sversion, char* msg, int32_t msgLen) {
930 931
  int32_t                 code;
  SStreamRecoverStep2Req* pReq = (SStreamRecoverStep2Req*)msg;
L
Liu Jicong 已提交
932
  SStreamTask*            pTask = streamMetaAcquireTask(pTq->pStreamMeta, pReq->taskId);
933 934 935 936 937
  if (pTask == NULL) {
    return -1;
  }

  // do recovery step 2
938
  code = streamSourceRecoverScanStep2(pTask, sversion);
939
  if (code < 0) {
L
Liu Jicong 已提交
940
    streamMetaReleaseTask(pTq->pStreamMeta, pTask);
941 942 943
    return -1;
  }

944
  if (atomic_load_8(&pTask->status.taskStatus) == TASK_STATUS__DROPPING) {
L
Liu Jicong 已提交
945 946 947 948
    streamMetaReleaseTask(pTq->pStreamMeta, pTask);
    return 0;
  }

949 950 951
  // restore param
  code = streamRestoreParam(pTask);
  if (code < 0) {
L
Liu Jicong 已提交
952
    streamMetaReleaseTask(pTq->pStreamMeta, pTask);
953 954 955 956 957 958
    return -1;
  }

  // set status normal
  code = streamSetStatusNormal(pTask);
  if (code < 0) {
L
Liu Jicong 已提交
959
    streamMetaReleaseTask(pTq->pStreamMeta, pTask);
960 961 962 963 964 965
    return -1;
  }

  // dispatch recover finish req to all related downstream task
  code = streamDispatchRecoverFinishReq(pTask);
  if (code < 0) {
L
Liu Jicong 已提交
966
    streamMetaReleaseTask(pTq->pStreamMeta, pTask);
967 968 969
    return -1;
  }

L
Liu Jicong 已提交
970 971 972
  atomic_store_8(&pTask->fillHistory, 0);
  streamMetaSaveTask(pTq->pStreamMeta, pTask);

L
Liu Jicong 已提交
973 974
  streamMetaReleaseTask(pTq->pStreamMeta, pTask);

975 976 977
  return 0;
}

L
Liu Jicong 已提交
978 979 980
int32_t tqProcessTaskRecoverFinishReq(STQ* pTq, SRpcMsg* pMsg) {
  char*   msg = POINTER_SHIFT(pMsg->pCont, sizeof(SMsgHead));
  int32_t msgLen = pMsg->contLen - sizeof(SMsgHead);
981 982

  // deserialize
983 984 985
  SStreamRecoverFinishReq req;

  SDecoder decoder;
X
Xiaoyu Wang 已提交
986
  tDecoderInit(&decoder, (uint8_t*)msg, msgLen);
987 988 989
  tDecodeSStreamRecoverFinishReq(&decoder, &req);
  tDecoderClear(&decoder);

990
  // find task
L
Liu Jicong 已提交
991
  SStreamTask* pTask = streamMetaAcquireTask(pTq->pStreamMeta, req.taskId);
992 993 994
  if (pTask == NULL) {
    return -1;
  }
995
  // do process request
996
  if (streamProcessRecoverFinishReq(pTask, req.childId) < 0) {
L
Liu Jicong 已提交
997
    streamMetaReleaseTask(pTq->pStreamMeta, pTask);
998 999 1000
    return -1;
  }

L
Liu Jicong 已提交
1001
  streamMetaReleaseTask(pTq->pStreamMeta, pTask);
1002
  return 0;
L
Liu Jicong 已提交
1003
}
L
Liu Jicong 已提交
1004

L
Liu Jicong 已提交
1005 1006 1007 1008 1009
int32_t tqProcessTaskRecoverFinishRsp(STQ* pTq, SRpcMsg* pMsg) {
  //
  return 0;
}

L
Liu Jicong 已提交
1010 1011 1012 1013 1014 1015 1016 1017 1018 1019 1020 1021 1022 1023 1024 1025
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 已提交
1026
  if (sz == 0 || pRes->affectedRows == 0) {
L
Liu Jicong 已提交
1027 1028 1029 1030 1031 1032 1033 1034 1035 1036 1037
    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);
1038
    colDataSetVal(pStartCol, i, (const char*)&pRes->skey, false);  // end key column
L
Liu Jicong 已提交
1039
    SColumnInfoData* pEndCol = taosArrayGet(pDelBlock->pDataBlock, END_TS_COLUMN_INDEX);
1040
    colDataSetVal(pEndCol, i, (const char*)&pRes->ekey, false);
L
Liu Jicong 已提交
1041 1042 1043
    // uid column
    SColumnInfoData* pUidCol = taosArrayGet(pDelBlock->pDataBlock, UID_COLUMN_INDEX);
    int64_t*         pUid = taosArrayGet(pRes->uidList, i);
1044
    colDataSetVal(pUidCol, i, (const char*)pUid, false);
L
Liu Jicong 已提交
1045

1046 1047 1048
    colDataSetNULL(taosArrayGet(pDelBlock->pDataBlock, GROUPID_COLUMN_INDEX), i);
    colDataSetNULL(taosArrayGet(pDelBlock->pDataBlock, CALCULATE_START_TS_COLUMN_INDEX), i);
    colDataSetNULL(taosArrayGet(pDelBlock->pDataBlock, CALCULATE_END_TS_COLUMN_INDEX), i);
L
Liu Jicong 已提交
1049 1050
  }

L
Liu Jicong 已提交
1051 1052
  taosArrayDestroy(pRes->uidList);

L
Liu Jicong 已提交
1053 1054 1055
  int32_t* pRef = taosMemoryMalloc(sizeof(int32_t));
  *pRef = 1;

L
Liu Jicong 已提交
1056 1057 1058 1059 1060 1061 1062
  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;

1063
    qDebug("delete req enqueue stream task: %d, ver: %" PRId64, pTask->id.taskId, ver);
L
Liu Jicong 已提交
1064

L
Liu Jicong 已提交
1065
    if (!failed) {
S
Shengliang Guan 已提交
1066
      SStreamRefDataBlock* pRefBlock = taosAllocateQitem(sizeof(SStreamRefDataBlock), DEF_QITEM, 0);
L
Liu Jicong 已提交
1067 1068 1069 1070 1071
      pRefBlock->type = STREAM_INPUT__REF_DATA_BLOCK;
      pRefBlock->pBlock = pDelBlock;
      pRefBlock->dataRef = pRef;
      atomic_add_fetch_32(pRefBlock->dataRef, 1);

1072
      if (tAppendDataToInputQueue(pTask, (SStreamQueueItem*)pRefBlock) < 0) {
1073
        qError("stream task input del failed, task id %d", pTask->id.taskId);
L
Liu Jicong 已提交
1074

L
Liu Jicong 已提交
1075
        atomic_sub_fetch_32(pRef, 1);
L
Liu Jicong 已提交
1076
        taosFreeQitem(pRefBlock);
L
Liu Jicong 已提交
1077 1078
        continue;
      }
L
Liu Jicong 已提交
1079

L
Liu Jicong 已提交
1080
      if (streamSchedExec(pTask) < 0) {
1081
        qError("stream task launch failed, task id %d", pTask->id.taskId);
L
Liu Jicong 已提交
1082 1083
        continue;
      }
L
Liu Jicong 已提交
1084

L
Liu Jicong 已提交
1085 1086 1087 1088
    } else {
      streamTaskInputFail(pTask);
    }
  }
L
Liu Jicong 已提交
1089

L
Liu Jicong 已提交
1090
  int32_t ref = atomic_sub_fetch_32(pRef, 1);
L
Liu Jicong 已提交
1091
  /*A(ref >= 0);*/
L
Liu Jicong 已提交
1092
  if (ref == 0) {
L
Liu Jicong 已提交
1093
    blockDataDestroy(pDelBlock);
L
Liu Jicong 已提交
1094 1095 1096 1097
    taosMemoryFree(pRef);
  }

#if 0
S
Shengliang Guan 已提交
1098
    SStreamDataBlock* pStreamBlock = taosAllocateQitem(sizeof(SStreamDataBlock), DEF_QITEM, 0);
L
Liu Jicong 已提交
1099 1100 1101 1102 1103 1104 1105 1106
    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) {
1107
      if (tAppendDataToInputQueue(pTask, (SStreamQueueItem*)pStreamBlock) < 0) {
1108
        qError("stream task input del failed, task id %d", pTask->id.taskId);
L
Liu Jicong 已提交
1109 1110 1111 1112
        continue;
      }

      if (streamSchedExec(pTask) < 0) {
1113
        qError("stream task launch failed, task id %d", pTask->id.taskId);
L
Liu Jicong 已提交
1114 1115 1116 1117 1118 1119
        continue;
      }
    } else {
      streamTaskInputFail(pTask);
    }
  }
L
Liu Jicong 已提交
1120
  blockDataDestroy(pDelBlock);
L
Liu Jicong 已提交
1121
#endif
L
Liu Jicong 已提交
1122 1123 1124 1125

  return 0;
}

1126 1127 1128
static int32_t addSubmitBlockNLaunchTask(STqOffsetStore* pOffsetStore, SStreamTask* pTask, SStreamDataSubmit2* pSubmit,
                                         const char* key, int64_t ver) {
  doSaveTaskOffset(pOffsetStore, key, ver);
1129
  int32_t code = tqAddInputBlockNLaunchTask(pTask, (SStreamQueueItem*)pSubmit, ver);
1130 1131 1132 1133 1134

  // remove the offset, if all functions are completed successfully.
  if (code == TSDB_CODE_SUCCESS) {
    tqOffsetDelete(pOffsetStore, key);
  }
1135 1136

  return code;
1137 1138
}

L
Liu Jicong 已提交
1139
int32_t tqProcessSubmitReq(STQ* pTq, SPackedData submit) {
1140
#if 0
1141
  void* pIter = NULL;
1142
  SStreamDataSubmit2* pSubmit = streamDataSubmitNew(submit, STREAM_INPUT__DATA_SUBMIT);
L
Liu Jicong 已提交
1143
  if (pSubmit == NULL) {
L
Liu Jicong 已提交
1144
    terrno = TSDB_CODE_OUT_OF_MEMORY;
L
Liu Jicong 已提交
1145
    tqError("failed to create data submit for stream since out of memory");
1146
    saveOffsetForAllTasks(pTq, submit.ver);
1147
    return -1;
L
Liu Jicong 已提交
1148 1149
  }

1150 1151
  SArray* pInputQueueFullTasks = taosArrayInit(4, POINTER_BYTES);

L
Liu Jicong 已提交
1152
  while (1) {
L
Liu Jicong 已提交
1153
    pIter = taosHashIterate(pTq->pStreamMeta->pTasks, pIter);
1154 1155 1156 1157
    if (pIter == NULL) {
      break;
    }

1158
    SStreamTask* pTask = *(SStreamTask**)pIter;
1159 1160 1161 1162
    if (pTask->taskLevel != TASK_LEVEL__SOURCE) {
      continue;
    }

1163
    if (pTask->status.taskStatus == TASK_STATUS__RECOVER_PREPARE || pTask->status.taskStatus == TASK_STATUS__WAIT_DOWNSTREAM) {
1164
      tqDebug("stream task:%d skip push data, not ready for processing, status %d", pTask->id.taskId,
1165
              pTask->status.taskStatus);
L
Liu Jicong 已提交
1166 1167
      continue;
    }
L
Liu Jicong 已提交
1168

1169 1170 1171
    // check if offset value exists
    char key[128] = {0};
    createStreamTaskOffsetKey(key, pTask->id.streamId, pTask->id.taskId);
1172

1173 1174 1175 1176 1177 1178 1179 1180
    if (tInputQueueIsFull(pTask)) {
      STqOffset* pOffset = tqOffsetRead(pTq->pOffsetStore, key);

      int64_t ver = submit.ver;
      if (pOffset == NULL) {
        doSaveTaskOffset(pTq->pOffsetStore, key, submit.ver);
      } else {
        ver = pOffset->val.version;
L
Liu Jicong 已提交
1181 1182
      }

1183 1184
      tqDebug("s-task:%s input queue is full, discard submit block, ver:%" PRId64, pTask->id.idStr, ver);
      taosArrayPush(pInputQueueFullTasks, &pTask);
1185 1186 1187 1188 1189
      continue;
    }

    // check if offset value exists
    STqOffset* pOffset = tqOffsetRead(pTq->pOffsetStore, key);
1190
    ASSERT(pOffset == NULL);
1191

1192
    addSubmitBlockNLaunchTask(pTq->pOffsetStore, pTask, pSubmit, key, submit.ver);
L
Liu Jicong 已提交
1193 1194
  }

1195 1196
  streamDataSubmitDestroy(pSubmit);
  taosFreeQitem(pSubmit);
1197
#endif
L
Liu Jicong 已提交
1198

1199
  tqStartStreamTasks(pTq);
1200
  return 0;
L
Liu Jicong 已提交
1201 1202
}

L
Liu Jicong 已提交
1203 1204
int32_t tqProcessTaskRunReq(STQ* pTq, SRpcMsg* pMsg) {
  SStreamTaskRunReq* pReq = pMsg->pCont;
1205 1206 1207 1208

  int32_t taskId = pReq->taskId;
  int32_t vgId = TD_VID(pTq->pVnode);

1209 1210
  if (taskId == WAL_READ_TASKS_ID) {  // all tasks are extracted submit data from the wal
    tqStreamTasksScanWal(pTq);
L
Liu Jicong 已提交
1211
    return 0;
1212
  }
1213

1214
  SStreamTask* pTask = streamMetaAcquireTask(pTq->pStreamMeta, taskId);
1215 1216 1217 1218 1219 1220 1221 1222
  if (pTask != NULL) {
    if (pTask->status.taskStatus == TASK_STATUS__NORMAL) {
      tqDebug("vgId:%d s-task:%s start to process run req", vgId, pTask->id.idStr);
      streamProcessRunReq(pTask);
    } else if (pTask->status.taskStatus == TASK_STATUS__RESTORE) {
      tqDebug("vgId:%d s-task:%s start to process block from wal, last chk point:%" PRId64, vgId,
              pTask->id.idStr, pTask->chkInfo.version);
      streamProcessRunReq(pTask);
1223
    } else {
1224
      tqDebug("vgId:%d s-task:%s ignore run req since not in ready state", vgId, pTask->id.idStr);
1225
    }
1226 1227 1228 1229 1230 1231 1232

    streamMetaReleaseTask(pTq->pStreamMeta, pTask);
    tqStartStreamTasks(pTq);
    return 0;
  } else {
    tqError("vgId:%d failed to found s-task, taskId:%d", vgId, taskId);
    return -1;
L
Liu Jicong 已提交
1233
  }
L
Liu Jicong 已提交
1234 1235
}

L
Liu Jicong 已提交
1236
int32_t tqProcessTaskDispatchReq(STQ* pTq, SRpcMsg* pMsg, bool exec) {
1237 1238 1239 1240 1241
  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 已提交
1242
  tDecoderInit(&decoder, (uint8_t*)msgBody, msgLen);
1243
  tDecodeStreamDispatchReq(&decoder, &req);
L
Liu Jicong 已提交
1244

1245
  SStreamTask* pTask = streamMetaAcquireTask(pTq->pStreamMeta, req.taskId);
L
Liu Jicong 已提交
1246
  if (pTask) {
1247
    SRpcMsg rsp = { .info = pMsg->info, .code = 0 };
L
Liu Jicong 已提交
1248
    streamProcessDispatchReq(pTask, &req, &rsp, exec);
L
Liu Jicong 已提交
1249
    streamMetaReleaseTask(pTq->pStreamMeta, pTask);
L
Liu Jicong 已提交
1250
    return 0;
1251 1252
  } else {
    return -1;
L
Liu Jicong 已提交
1253
  }
L
Liu Jicong 已提交
1254 1255
}

L
Liu Jicong 已提交
1256 1257
int32_t tqProcessTaskDispatchRsp(STQ* pTq, SRpcMsg* pMsg) {
  SStreamDispatchRsp* pRsp = POINTER_SHIFT(pMsg->pCont, sizeof(SMsgHead));
1258
  int32_t             taskId = ntohl(pRsp->upstreamTaskId);
L
Liu Jicong 已提交
1259
  SStreamTask*        pTask = streamMetaAcquireTask(pTq->pStreamMeta, taskId);
1260
  tqDebug("recv dispatch rsp, code:%x", pMsg->code);
L
Liu Jicong 已提交
1261
  if (pTask) {
1262
    streamProcessDispatchRsp(pTask, pRsp, pMsg->code);
L
Liu Jicong 已提交
1263
    streamMetaReleaseTask(pTq->pStreamMeta, pTask);
L
Liu Jicong 已提交
1264
    return 0;
1265 1266
  } else {
    return -1;
L
Liu Jicong 已提交
1267
  }
L
Liu Jicong 已提交
1268
}
L
Liu Jicong 已提交
1269

1270
int32_t tqProcessTaskDropReq(STQ* pTq, int64_t sversion, char* msg, int32_t msgLen) {
L
Liu Jicong 已提交
1271
  SVDropStreamTaskReq* pReq = (SVDropStreamTaskReq*)msg;
1272
  streamMetaRemoveTask(pTq->pStreamMeta, pReq->taskId);
L
Liu Jicong 已提交
1273
  return 0;
L
Liu Jicong 已提交
1274
}
L
Liu Jicong 已提交
1275 1276 1277 1278 1279 1280 1281

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;
1282
  tDecoderInit(&decoder, (uint8_t*)msgBody, msgLen);
L
Liu Jicong 已提交
1283
  tDecodeStreamRetrieveReq(&decoder, &req);
L
Liu Jicong 已提交
1284
  tDecoderClear(&decoder);
L
Liu Jicong 已提交
1285
  int32_t      taskId = req.dstTaskId;
L
Liu Jicong 已提交
1286
  SStreamTask* pTask = streamMetaAcquireTask(pTq->pStreamMeta, taskId);
L
Liu Jicong 已提交
1287
  if (pTask) {
1288
    SRpcMsg rsp = { .info = pMsg->info, .code = 0 };
L
Liu Jicong 已提交
1289
    streamProcessRetrieveReq(pTask, &req, &rsp);
L
Liu Jicong 已提交
1290
    streamMetaReleaseTask(pTq->pStreamMeta, pTask);
L
Liu Jicong 已提交
1291
    tDeleteStreamRetrieveReq(&req);
L
Liu Jicong 已提交
1292
    return 0;
L
Liu Jicong 已提交
1293 1294
  } else {
    return -1;
L
Liu Jicong 已提交
1295 1296 1297 1298 1299 1300 1301
  }
}

int32_t tqProcessTaskRetrieveRsp(STQ* pTq, SRpcMsg* pMsg) {
  //
  return 0;
}
L
Liu Jicong 已提交
1302

1303 1304 1305 1306 1307 1308
int32_t vnodeEnqueueStreamMsg(SVnode* pVnode, SRpcMsg* pMsg) {
  STQ*      pTq = pVnode->pTq;
  SMsgHead* msgStr = pMsg->pCont;
  char*     msgBody = POINTER_SHIFT(msgStr, sizeof(SMsgHead));
  int32_t   msgLen = pMsg->contLen - sizeof(SMsgHead);
  int32_t   code = 0;
L
Liu Jicong 已提交
1309 1310 1311

  SStreamDispatchReq req;
  SDecoder           decoder;
1312
  tDecoderInit(&decoder, (uint8_t*)msgBody, msgLen);
L
Liu Jicong 已提交
1313 1314
  if (tDecodeStreamDispatchReq(&decoder, &req) < 0) {
    code = TSDB_CODE_MSG_DECODE_ERROR;
L
Liu Jicong 已提交
1315
    tDecoderClear(&decoder);
L
Liu Jicong 已提交
1316 1317
    goto FAIL;
  }
L
Liu Jicong 已提交
1318
  tDecoderClear(&decoder);
L
Liu Jicong 已提交
1319

L
Liu Jicong 已提交
1320
  int32_t taskId = req.taskId;
L
Liu Jicong 已提交
1321

L
Liu Jicong 已提交
1322
  SStreamTask* pTask = streamMetaAcquireTask(pTq->pStreamMeta, taskId);
L
Liu Jicong 已提交
1323
  if (pTask) {
1324
    SRpcMsg rsp = { .info = pMsg->info, .code = 0 };
L
Liu Jicong 已提交
1325
    streamProcessDispatchReq(pTask, &req, &rsp, false);
L
Liu Jicong 已提交
1326
    streamMetaReleaseTask(pTq->pStreamMeta, pTask);
L
Liu Jicong 已提交
1327 1328
    rpcFreeCont(pMsg->pCont);
    taosFreeQitem(pMsg);
1329
    return 0;
5
54liuyao 已提交
1330 1331
  } else {
    tDeleteStreamDispatchReq(&req);
L
Liu Jicong 已提交
1332
  }
L
Liu Jicong 已提交
1333

1334 1335
  code = TSDB_CODE_STREAM_TASK_NOT_EXIST;

L
Liu Jicong 已提交
1336
FAIL:
1337 1338 1339 1340
  if (pMsg->info.handle == NULL) return -1;

  SMsgHead* pRspHead = rpcMallocCont(sizeof(SMsgHead) + sizeof(SStreamDispatchRsp));
  if (pRspHead == NULL) {
1341
    SRpcMsg rsp = { .code = TSDB_CODE_OUT_OF_MEMORY, .info = pMsg->info };
1342 1343 1344 1345 1346 1347 1348 1349 1350 1351 1352 1353 1354 1355 1356 1357
    tqDebug("send dispatch error rsp, code: %x", code);
    tmsgSendRsp(&rsp);
    rpcFreeCont(pMsg->pCont);
    taosFreeQitem(pMsg);
    return -1;
  }

  pRspHead->vgId = htonl(req.upstreamNodeId);
  SStreamDispatchRsp* pRsp = POINTER_SHIFT(pRspHead, sizeof(SMsgHead));
  pRsp->streamId = htobe64(req.streamId);
  pRsp->upstreamTaskId = htonl(req.upstreamTaskId);
  pRsp->upstreamNodeId = htonl(req.upstreamNodeId);
  pRsp->downstreamNodeId = htonl(pVnode->config.vgId);
  pRsp->downstreamTaskId = htonl(req.taskId);
  pRsp->inputStatus = TASK_OUTPUT_STATUS__NORMAL;

L
Liu Jicong 已提交
1358
  SRpcMsg rsp = {
1359
      .code = code, .info = pMsg->info, .contLen = sizeof(SMsgHead) + sizeof(SStreamDispatchRsp), .pCont = pRspHead};
1360
  tqDebug("send dispatch error rsp, code: %x", code);
L
Liu Jicong 已提交
1361
  tmsgSendRsp(&rsp);
L
Liu Jicong 已提交
1362 1363
  rpcFreeCont(pMsg->pCont);
  taosFreeQitem(pMsg);
1364
  return -1;
L
Liu Jicong 已提交
1365
}
L
Liu Jicong 已提交
1366

1367
int32_t tqCheckLogInWal(STQ* pTq, int64_t sversion) { return sversion <= pTq->walLogLastVer; }
1368

1369
int32_t tqStartStreamTasks(STQ* pTq) {
1370 1371
  int32_t vgId = TD_VID(pTq->pVnode);

1372 1373
  SStreamMeta* pMeta = pTq->pStreamMeta;
  taosWLockLatch(&pMeta->lock);
1374 1375 1376 1377 1378 1379 1380
  int32_t numOfTasks = taosHashGetSize(pTq->pStreamMeta->pTasks);
  if (numOfTasks == 0) {
    tqInfo("vgId:%d no stream tasks exists", vgId);
    taosWUnLockLatch(&pTq->pStreamMeta->lock);
    return 0;
  }

1381 1382 1383 1384 1385 1386 1387 1388
  pMeta->walScan += 1;

  if (pMeta->walScan > 1) {
    tqDebug("vgId:%d wal read task has been launched, remain scan times:%d", vgId, pMeta->walScan);
    taosWUnLockLatch(&pTq->pStreamMeta->lock);
    return 0;
  }

1389 1390 1391 1392
  SStreamTaskRunReq* pRunReq = rpcMallocCont(sizeof(SStreamTaskRunReq));
  if (pRunReq == NULL) {
    terrno = TSDB_CODE_OUT_OF_MEMORY;
    tqError("vgId:%d failed restore stream tasks, code:%s", vgId, terrstr(terrno));
1393
    taosWUnLockLatch(&pTq->pStreamMeta->lock);
1394 1395 1396
    return -1;
  }

1397
  tqInfo("vgId:%d start wal scan stream tasks, tasks:%d", vgId, numOfTasks);
1398 1399
  initOffsetForAllRestoreTasks(pTq);

1400 1401
  pRunReq->head.vgId = vgId;
  pRunReq->streamId = 0;
1402
  pRunReq->taskId = WAL_READ_TASKS_ID;
1403 1404 1405

  SRpcMsg msg = {.msgType = TDMT_STREAM_TASK_RUN, .pCont = pRunReq, .contLen = sizeof(SStreamTaskRunReq)};
  tmsgPutToQueue(&pTq->pVnode->msgCb, STREAM_QUEUE, &msg);
1406
  taosWUnLockLatch(&pTq->pStreamMeta->lock);
1407 1408 1409

  return 0;
}