• Home
  • Features
  • Pricing
  • Docs
  • Announcements
  • Sign In

taosdata / TDengine / #5034

24 Apr 2026 11:25AM UTC coverage: 73.058%. Remained the same
#5034

push

travis-ci

web-flow
merge: from main to 3.0 branch #35224

merge: from main to 3.0 branch[manual-only]

1336 of 1975 new or added lines in 48 files covered. (67.65%)

14149 existing lines in 164 files now uncovered.

275896 of 377640 relevant lines covered (73.06%)

132944440.29 hits per line

Source File
Press 'n' to go to next uncovered line, 'b' for previous

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

16
#include "catalog.h"
17
#include "command.h"
18
#include "query.h"
19
#include "schInt.h"
20
#include "tarray.h"
21
#include "tglobal.h"
22
#include "tmisce.h"
23
#include "tmsg.h"
24
#include "tref.h"
25
#include "trpc.h"
26

27
// clang-format off
28
int32_t schValidateRspMsgType(SSchJob *pJob, SSchTask *pTask, int32_t msgType) {
1,638,196,436✔
29
  int32_t lastMsgType = pTask->lastMsgType;
1,638,196,436✔
30
  int32_t taskStatus = SCH_GET_TASK_STATUS(pTask);
1,638,217,863✔
31
  int32_t reqMsgType = (msgType & 1U) ? msgType : (msgType - 1);
1,638,222,435✔
32
  switch (msgType) {
1,638,222,435✔
33
    case TDMT_SCH_LINK_BROKEN:
86,177,352✔
34
    case TDMT_SCH_EXPLAIN_RSP:
35
      return TSDB_CODE_SUCCESS;
86,177,352✔
36
    case TDMT_SCH_FETCH_RSP:
280,595,319✔
37
    case TDMT_SCH_MERGE_FETCH_RSP:
38
      if (lastMsgType != reqMsgType) {
280,595,319✔
39
        SCH_TASK_ELOG("rsp msg type mis-match, last sent msgType:%s, rspType:%s", TMSG_INFO(lastMsgType),
2,020✔
40
                      TMSG_INFO(msgType));
41
        SCH_ERR_RET(TSDB_CODE_QW_MSG_ERROR);
2,020✔
42
      }
43
      if (taskStatus != JOB_TASK_STATUS_FETCH) {
280,593,299✔
44
        SCH_TASK_ELOG("rsp msg conflicted with task status, status:%s, rspType:%s", jobTaskStatusStr(taskStatus),
2,020✔
45
                      TMSG_INFO(msgType));
46
        SCH_ERR_RET(TSDB_CODE_QW_MSG_ERROR);
2,020✔
47
      }
48

49
      return TSDB_CODE_SUCCESS;
280,591,243✔
50
    case TDMT_SCH_MERGE_QUERY_RSP:
1,271,425,467✔
51
    case TDMT_SCH_QUERY_RSP:
52
    case TDMT_VND_CREATE_TABLE_RSP:
53
    case TDMT_VND_DROP_TABLE_RSP:
54
    case TDMT_VND_ALTER_TABLE_RSP:
55
    case TDMT_VND_SUBMIT_RSP:
56
    case TDMT_VND_DELETE_RSP:
57
    case TDMT_VND_COMMIT_RSP:
58
      break;
1,271,425,467✔
59
    default:
24,297✔
60
      SCH_TASK_ELOG("unknown rsp msg, type:%s, status:%s", TMSG_INFO(msgType), jobTaskStatusStr(taskStatus));
24,297✔
61
      SCH_ERR_RET(TSDB_CODE_INVALID_MSG);
24,297✔
62
  }
63

64
  if (lastMsgType != reqMsgType) {
1,271,447,744✔
65
    SCH_TASK_ELOG("rsp msg type mis-match, last sent msgType:%s, rspType:%s", TMSG_INFO(lastMsgType),
2,020✔
66
                  TMSG_INFO(msgType));
67
    SCH_ERR_RET(TSDB_CODE_QW_MSG_ERROR);
2,020✔
68
  }
69

70
  if (taskStatus != JOB_TASK_STATUS_EXEC) {
1,271,445,724✔
71
    SCH_TASK_ELOG("rsp msg conflicted with task status, status:%s, rspType:%s", jobTaskStatusStr(taskStatus),
2,020✔
72
                  TMSG_INFO(msgType));
73
    SCH_ERR_RET(TSDB_CODE_QW_MSG_ERROR);
972✔
74
  }
75

76
  return TSDB_CODE_SUCCESS;
1,271,419,994✔
77
}
78

79
int32_t schProcessFetchRsp(SSchJob *pJob, SSchTask *pTask, char *msg, int32_t rspCode) {
280,600,121✔
80
  SRetrieveTableRsp *rsp = (SRetrieveTableRsp *)msg;
280,600,121✔
81
  int32_t code = 0;
280,600,121✔
82
  
83
  SCH_ERR_JRET(rspCode);
280,600,121✔
84

85
  if (NULL == msg) {
280,598,101✔
86
    SCH_ERR_RET(TSDB_CODE_QRY_INVALID_INPUT);
2,020✔
87
  }
88

89

90
  if (SCH_IS_EXPLAIN_JOB(pJob)) {
280,596,081✔
91
    if (rsp->completed) {
2,193,238✔
92
      SRetrieveTableRsp *pRsp = NULL;
1,117,141✔
93
      SCH_ERR_JRET(qExecExplainEnd(SCH_JOB_EXPLAIN_CTX(pJob), &pRsp));
1,117,141✔
94
      if (pRsp) {
1,117,141✔
95
        SCH_ERR_JRET(schProcessOnExplainDone(SCH_PARENT_JOB(pJob), pTask, pRsp));
1,112,975✔
96
      } else {
97
        SCH_ERR_JRET(schNotifyJobAllTasks(SCH_PARENT_JOB(pJob), pTask, TASK_NOTIFY_FINISHED));
4,166✔
98
      }
99
  
100
      taosMemoryFreeClear(msg);
1,117,141✔
101
  
102
      return TSDB_CODE_SUCCESS;
1,117,141✔
103
    }
104
  
105
    SCH_ERR_JRET(schLaunchFetchTask(pJob));
1,076,097✔
106
  
107
    taosMemoryFreeClear(msg);
1,076,097✔
108
  
109
    return TSDB_CODE_SUCCESS;
1,076,097✔
110
  }
111
  
112
  if (pJob->fetchRes) {
278,402,783✔
113
    SCH_TASK_ELOG("got fetch rsp while res already exists, res:%p", pJob->fetchRes);
2,020✔
114
    SCH_ERR_JRET(TSDB_CODE_SCH_STATUS_ERROR);
2,020✔
115
  }
116
  
117
  atomic_store_ptr(&pJob->fetchRes, rsp);
278,400,831✔
118
  (void)atomic_add_fetch_64(&pJob->resNumOfRows, htobe64(rsp->numOfRows));
278,400,834✔
119
  
120
  if (rsp->completed) {
278,400,870✔
121
    SCH_SET_TASK_STATUS(pTask, JOB_TASK_STATUS_SUCC);
143,580,127✔
122
  }
123
  
124
  SCH_TASK_DLOG("got fetch rsp, rows:%" PRId64 ", complete:%d", htobe64(rsp->numOfRows), rsp->completed);
278,400,870✔
125

126
  msg = NULL;
278,400,870✔
127
  schProcessOnDataFetched(pJob);
278,400,870✔
128

129
_return:
278,404,873✔
130

131
  taosMemoryFreeClear(msg);
278,404,873✔
132

133
  SCH_RET(code);
278,404,873✔
134
}
135

136
int32_t schProcessExplainRsp(SSchJob *pJob, SSchTask *pTask, SExplainRsp *rsp) {
86,177,199✔
137
  SRetrieveTableRsp *pRsp = NULL;
86,177,199✔
138
  SExplainCtx* pCtx = SCH_JOB_EXPLAIN_CTX(pJob);
86,177,614✔
139
  SCH_ERR_RET(qExplainUpdateExecInfo(pCtx, qExplainGetCurrPlan(pCtx, pJob->subJobId), rsp, pTask->plan->id.groupId, &pRsp));
86,177,614✔
140
  
141
  if (pRsp) {
86,184,403✔
142
    SCH_ERR_RET(schProcessOnExplainDone(SCH_PARENT_JOB(pJob), pTask, pRsp));
4,722✔
143
  }
144

145
  return TSDB_CODE_SUCCESS;
86,184,403✔
146
}
147

148
int32_t schProcessResponseMsg(SSchJob *pJob, SSchTask *pTask, SDataBuf *pMsg, int32_t rspCode) {
1,633,507,699✔
149
  int32_t code = 0;
1,633,507,699✔
150
  int32_t msgSize = pMsg->len;
1,633,507,699✔
151
  int32_t msgType = pMsg->msgType;
1,633,515,737✔
152

153
  switch (msgType) {
1,633,485,788✔
154
    case TDMT_VND_COMMIT_RSP: {
4,033,682✔
155
      SCH_ERR_JRET(rspCode);
4,033,682✔
156
      SCH_ERR_JRET(schProcessOnTaskSuccess(pJob, pTask));
4,032,929✔
157
      break;
4,032,176✔
158
    }
159
    case TDMT_VND_CREATE_TABLE_RSP: {
47,997,782✔
160
      SVCreateTbBatchRsp batchRsp = {0};
47,997,782✔
161
      if (pMsg->pData) {
48,000,035✔
162
        SDecoder coder = {0};
47,993,504✔
163
        tDecoderInit(&coder, pMsg->pData, msgSize);
47,993,979✔
164
        code = tDecodeSVCreateTbBatchRsp(&coder, &batchRsp);
47,993,286✔
165
        if (TSDB_CODE_SUCCESS == code && batchRsp.nRsps > 0) {
47,987,584✔
166
          SCH_LOCK(SCH_WRITE, &pJob->resLock);
47,988,348✔
167
          if (NULL == pJob->execRes.res) {
47,986,497✔
168
            pJob->execRes.res = (void*)taosArrayInit(batchRsp.nRsps, POINTER_BYTES);
46,365,933✔
169
            if (NULL == pJob->execRes.res) {
46,370,646✔
170
              code = terrno;
×
171
              SCH_UNLOCK(SCH_WRITE, &pJob->resLock);
×
172
              
173
              tDecoderClear(&coder);
×
174
              SCH_ERR_JRET(code);
×
175
            }
176
            
177
            pJob->execRes.msgType = TDMT_VND_CREATE_TABLE;
46,370,527✔
178
          }
179

180
          for (int32_t i = 0; i < batchRsp.nRsps; ++i) {
101,983,735✔
181
            SVCreateTbRsp *rsp = batchRsp.pRsps + i;
53,993,515✔
182
            if (rsp->pMeta) {
53,989,820✔
183
              if (NULL == taosArrayPush((SArray*)pJob->execRes.res, &rsp->pMeta)) {
107,943,956✔
184
                code = terrno;
×
185
                SCH_UNLOCK(SCH_WRITE, &pJob->resLock);
×
186
                
187
                tDecoderClear(&coder);
×
188
                SCH_ERR_JRET(code);
×
189
              }
190
            }
191
            
192
            if (TSDB_CODE_SUCCESS != rsp->code) {
53,993,562✔
193
              code = rsp->code;
10,392✔
194
            }
195
          }
196

197
          if (taosArrayGetSize((SArray*)pJob->execRes.res) <= 0) {        
47,990,220✔
198
            taosArrayDestroy((SArray*)pJob->execRes.res);
18,192✔
199
            pJob->execRes.res = NULL;
18,192✔
200
          }
201
          SCH_UNLOCK(SCH_WRITE, &pJob->resLock);
47,990,256✔
202
        }
203
        
204
        tDecoderClear(&coder);
47,990,968✔
205
        SCH_ERR_JRET(code);
47,990,734✔
206
      }
207

208
      SCH_ERR_JRET(rspCode);
47,978,667✔
209
      taosMemoryFreeClear(pMsg->pData);
47,971,539✔
210

211
      SCH_ERR_JRET(schProcessOnTaskSuccess(pJob, pTask));
47,978,579✔
212
      break;
47,978,036✔
213
    }
214
    case TDMT_VND_DROP_TABLE_RSP: {
1,845,796✔
215
      SVDropTbBatchRsp batchRsp = {0};
1,845,796✔
216
      if (pMsg->pData) {
1,845,796✔
217
        SDecoder coder = {0};
1,845,188✔
218
        tDecoderInit(&coder, pMsg->pData, msgSize);
1,845,188✔
219
        code = tDecodeSVDropTbBatchRsp(&coder, &batchRsp);
1,845,188✔
220
        if (TSDB_CODE_SUCCESS == code && batchRsp.nRsps > 0) {
1,845,188✔
221
          for (int32_t i = 0; i < batchRsp.nRsps; ++i) {
3,748,820✔
222
            SVDropTbRsp *rsp = batchRsp.pRsps + i;
1,903,632✔
223
            if (TSDB_CODE_SUCCESS != rsp->code) {
1,903,314✔
224
              code = rsp->code;
×
225
              tDecoderClear(&coder);
×
226
              SCH_ERR_JRET(code);
×
227
            }
228
          }
229
        }
230
        tDecoderClear(&coder);
1,845,188✔
231
        SCH_ERR_JRET(code);
1,845,188✔
232
      }
233

234
      SCH_ERR_JRET(rspCode);
1,845,796✔
235
      taosMemoryFreeClear(pMsg->pData);
1,845,188✔
236

237
      SCH_ERR_JRET(schProcessOnTaskSuccess(pJob, pTask));
1,845,188✔
238
      break;
1,845,188✔
239
    }
240
    case TDMT_VND_ALTER_TABLE_RSP: {
16,226,091✔
241
      SVAlterTbRsp rsp = {0};
16,226,091✔
242
      if (pMsg->pData) {
16,226,091✔
243
        SDecoder coder = {0};
16,223,752✔
244
        tDecoderInit(&coder, pMsg->pData, msgSize);
16,223,752✔
245
        code = tDecodeSVAlterTbRsp(&coder, &rsp);
16,223,752✔
246
        tDecoderClear(&coder);
16,223,752✔
247
        SCH_ERR_JRET(code);
16,223,752✔
248
        if (rsp.code == TSDB_CODE_VND_SAME_TAG) {
16,223,752✔
249
          rsp.code = TSDB_CODE_SUCCESS;
×
250
        }
251
        SCH_ERR_JRET(rsp.code);
16,223,752✔
252

253
        pJob->execRes.res = rsp.pMeta;
15,629,016✔
254
        pJob->execRes.msgType = TDMT_VND_ALTER_TABLE;
15,629,016✔
255
      }
256

257
      SCH_ERR_JRET(rspCode);
15,631,355✔
258

259
      if (NULL == pMsg->pData) {
15,631,036✔
260
        SCH_ERR_JRET(TSDB_CODE_QRY_INVALID_INPUT);
2,020✔
261
      }
262

263
      taosMemoryFreeClear(pMsg->pData);
15,629,016✔
264

265
      SCH_ERR_JRET(schProcessOnTaskSuccess(pJob, pTask));
15,629,016✔
266
      break;
15,626,911✔
267
    }
268
    case TDMT_VND_SUBMIT_RSP: {
649,588,549✔
269
      SCH_ERR_JRET(rspCode);
649,588,549✔
270

271
      if (pMsg->pData) {
649,491,635✔
272
        SDecoder    coder = {0};
649,498,990✔
273
        SSubmitRsp2 *rsp = taosMemoryMalloc(sizeof(*rsp));
649,503,269✔
274
        if (NULL == rsp) {
649,495,812✔
275
          SCH_ERR_JRET(terrno);
2,020✔
276
        }
277
        tDecoderInit(&coder, pMsg->pData, msgSize);
649,495,812✔
278
        code = tDecodeSSubmitRsp2(&coder, rsp);
649,498,904✔
279
        tDecoderClear(&coder);
649,481,657✔
280
        if (code) {
649,489,178✔
281
          SCH_TASK_ELOG("tDecodeSSubmitRsp2 failed, code:%d", code);
2,020✔
282
          tDestroySSubmitRsp2(rsp, TSDB_MSG_FLG_DECODE);
2,020✔
283
          taosMemoryFree(rsp);
2,020✔
284
          SCH_ERR_JRET(code);
2,020✔
285
        }
286

287
        (void)atomic_add_fetch_64(&pJob->resNumOfRows, rsp->affectedRows);
649,487,158✔
288

289
        int32_t createTbRspNum = taosArrayGetSize(rsp->aCreateTbRsp);
649,507,815✔
290
        SCH_TASK_DLOG("submit succeed, affectedRows:%d, createTbRspNum:%d", rsp->affectedRows, createTbRspNum);
649,489,167✔
291

292
        if (rsp->aCreateTbRsp && taosArrayGetSize(rsp->aCreateTbRsp) > 0) {
649,494,217✔
293
          SCH_LOCK(SCH_WRITE, &pJob->resLock);
13,645,665✔
294
          if (pJob->execRes.res) {
13,642,749✔
295
            SSubmitRsp2 *sum = pJob->execRes.res;
92,060✔
296
            sum->affectedRows += rsp->affectedRows;
92,060✔
297
            if (sum->aCreateTbRsp) {
92,060✔
298
              if (NULL == taosArrayAddAll(sum->aCreateTbRsp, rsp->aCreateTbRsp)) {
90,679✔
299
                code = terrno;
×
300
                SCH_UNLOCK(SCH_WRITE, &pJob->resLock);
×
301
                SCH_ERR_JRET(code);
×
302
              }
303
              
304
              taosArrayDestroy(rsp->aCreateTbRsp);
90,679✔
305
            } else {
306
              TSWAP(sum->aCreateTbRsp, rsp->aCreateTbRsp);
1,381✔
307
            }
308
            taosMemoryFree(rsp);
92,060✔
309
          } else {
310
            pJob->execRes.res = rsp;
13,549,603✔
311
            pJob->execRes.msgType = TDMT_VND_SUBMIT;
13,552,861✔
312
          }
313
          pJob->execRes.numOfBytes += pTask->msgLen;
13,640,190✔
314
          SCH_UNLOCK(SCH_WRITE, &pJob->resLock);
13,643,820✔
315
        } else {
316
          SCH_LOCK(SCH_WRITE, &pJob->resLock);
635,849,458✔
317
          pJob->execRes.numOfBytes += pTask->msgLen;
635,846,734✔
318
          if (NULL == pJob->execRes.res) {
635,850,196✔
319
            TSWAP(pJob->execRes.res, rsp);
619,834,777✔
320
            pJob->execRes.msgType = TDMT_VND_SUBMIT;
619,819,282✔
321
          }
322
          SCH_UNLOCK(SCH_WRITE, &pJob->resLock);
635,849,658✔
323
          tDestroySSubmitRsp2(rsp, TSDB_MSG_FLG_DECODE);
635,850,884✔
324
          taosMemoryFree(rsp);
635,839,138✔
325
        }
326
      }
327

328
      taosMemoryFreeClear(pMsg->pData);
649,498,572✔
329

330
      SCH_ERR_JRET(schProcessOnTaskSuccess(pJob, pTask));
649,490,784✔
331

332
      break;
649,481,455✔
333
    }
334
    case TDMT_VND_DELETE_RSP: {
1,927,732✔
335
      SCH_ERR_JRET(rspCode);
1,927,732✔
336

337
      if (pMsg->pData) {
1,927,732✔
338
        SDecoder    coder = {0};
1,927,732✔
339
        SVDeleteRsp rsp = {0};
1,927,732✔
340
        tDecoderInit(&coder, pMsg->pData, msgSize);
1,927,732✔
341
        if (tDecodeSVDeleteRsp(&coder, &rsp) < 0) {
1,927,732✔
342
          code = terrno;
2,020✔
343
          tDecoderClear(&coder);
2,020✔
344
          SCH_ERR_JRET(code);
2,020✔
345
        }
346
        tDecoderClear(&coder);
1,925,698✔
347

348
        (void)atomic_add_fetch_64(&pJob->resNumOfRows, rsp.affectedRows);
1,925,712✔
349
        SCH_TASK_DLOG("delete succeed, affectedRows:%" PRId64, rsp.affectedRows);
1,925,712✔
350
      }
351

352
      taosMemoryFreeClear(pMsg->pData);
1,925,712✔
353

354
      SCH_ERR_JRET(schProcessOnTaskSuccess(pJob, pTask));
1,925,712✔
355

356
      break;
1,925,712✔
357
    }
358
    case TDMT_SCH_QUERY_RSP:
545,100,809✔
359
    case TDMT_SCH_MERGE_QUERY_RSP: {
360
      SCH_ERR_JRET(rspCode);
545,104,849✔
361
      if (NULL == pMsg->pData) {
531,417,118✔
362
        SCH_ERR_JRET(TSDB_CODE_QRY_INVALID_INPUT);
2,020✔
363
      }
364

365
      if (taosArrayGetSize(pTask->parents) == 0 && SCH_IS_EXPLAIN_JOB(pJob) && SCH_IS_INSERT_JOB(pJob)) {
531,421,215✔
366
        SRetrieveTableRsp *pRsp = NULL;
5,004✔
367
        SCH_ERR_JRET(qExecExplainEnd(SCH_JOB_EXPLAIN_CTX(pJob), &pRsp));
5,004✔
368
        if (pRsp) {
5,004✔
369
          SCH_ERR_JRET(schProcessOnExplainDone(SCH_PARENT_JOB(pJob), pTask, pRsp));
4,448✔
370
        }
371
      }
372

373
      SQueryTableRsp rsp = {0};
531,433,983✔
374
      if (tDeserializeSQueryTableRsp(pMsg->pData, msgSize, &rsp) < 0) {
531,436,628✔
375
        SCH_TASK_ELOG("tDeserializeSQueryTableRsp failed, msgSize:%d", msgSize);
2,020✔
376
        SCH_ERR_JRET(TSDB_CODE_QRY_INVALID_MSG);
541✔
377
      }
378
      
379
      SCH_ERR_JRET(rsp.code);
531,415,556✔
380

381
      SCH_ERR_JRET(schSaveJobExecRes(pJob, &rsp));
531,415,556✔
382

383
      (void)atomic_add_fetch_64(&pJob->resNumOfRows, rsp.affectedRows);
531,448,509✔
384

385
      taosMemoryFreeClear(pMsg->pData);
531,454,151✔
386

387
      SCH_ERR_JRET(schProcessOnTaskSuccess(pJob, pTask));
531,451,422✔
388

389
      break;
531,438,351✔
390
    }
391
    case TDMT_SCH_EXPLAIN_RSP: {
86,173,870✔
392
      SCH_ERR_JRET(rspCode);
86,181,950✔
393
      if (NULL == pMsg->pData) {
86,173,870✔
394
        SCH_ERR_JRET(TSDB_CODE_QRY_INVALID_INPUT);
2,020✔
395
      }
396

397
      if (!SCH_IS_EXPLAIN_JOB(pJob)) {
86,172,680✔
398
        SCH_TASK_ELOG("invalid msg received for none explain query, msg type:%s", TMSG_INFO(msgType));
2,020✔
399
        SCH_ERR_JRET(TSDB_CODE_QRY_INVALID_INPUT);
2,020✔
400
      }
401

402
      if (pJob->fetchRes) {
86,170,245✔
403
        SCH_TASK_ELOG("explain result is already generated, res:%p", pJob->fetchRes);
2,020✔
404
        SCH_ERR_JRET(TSDB_CODE_SCH_STATUS_ERROR);
2,020✔
405
      }
406

407
      SExplainRsp rsp = {0};
86,168,225✔
408
      if (tDeserializeSExplainRsp(pMsg->pData, msgSize, &rsp)) {
86,168,225✔
409
        tFreeSExplainRsp(&rsp);
2,020✔
410
        SCH_ERR_JRET(TSDB_CODE_OUT_OF_MEMORY);
2,020✔
411
      }
412

413
      SCH_ERR_JRET(schProcessExplainRsp(pJob, pTask, &rsp));
86,179,618✔
414

415
      taosMemoryFreeClear(pMsg->pData);
86,184,045✔
416
      break;
86,187,239✔
417
    }
418
    case TDMT_SCH_FETCH_RSP:
280,591,227✔
419
    case TDMT_SCH_MERGE_FETCH_RSP: {
420
      code = schProcessFetchRsp(pJob, pTask, pMsg->pData, rspCode);
280,591,227✔
421
      pMsg->pData = NULL;
280,591,234✔
422
      SCH_ERR_JRET(code);
280,591,242✔
423
      break;
280,591,242✔
424
    }
425
    case TDMT_SCH_DROP_TASK_RSP: {
2,020✔
426
      // NEVER REACH HERE
427
      SCH_TASK_ELOG("invalid status to handle drop task rsp, refId:0x%" PRIx64, pJob->refId);
2,020✔
428
      SCH_ERR_JRET(TSDB_CODE_SCH_INTERNAL_ERROR);
2,020✔
429
      break;
×
430
    }
431
    case TDMT_SCH_LINK_BROKEN:
2,020✔
432
      SCH_TASK_ELOG("link broken received, error:%x - %s", rspCode, tstrerror(rspCode));
2,020✔
433
      SCH_ERR_JRET(rspCode);
2,020✔
434
      break;
2,020✔
435
    case TDMT_MND_CREATE_TOKEN_RSP:{
×
436
      SCreateTokenRsp batchRsp = {0};
×
437
      code = tDeserializeSCreateTokenResp(pMsg->pData, msgSize, &batchRsp);
×
438
      SCH_ERR_JRET(code);
×
439
      break;
×
440
    }
441
    case TDMT_MND_CREATE_TOTP_SECRET:{
×
442
      SCreateTotpSecretRsp rsp = {0};
×
443
      code = tDeserializeSCreateTotpSecretRsp(pMsg->pData, msgSize, &rsp);
×
444
      SCH_ERR_JRET(code);
×
445
      break;
×
446
    }
447
    default:
125✔
448
      SCH_TASK_ELOG("unknown rsp msg, type:%d, status:%s", msgType, SCH_GET_TASK_STATUS_STR(pTask));
125✔
449
      SCH_ERR_JRET(TSDB_CODE_QRY_INVALID_INPUT);
137✔
450
  }
451

452
  return TSDB_CODE_SUCCESS;
1,619,099,603✔
453

454
_return:
14,424,986✔
455

456
  taosMemoryFreeClear(pMsg->pData);
14,424,986✔
457

458
  SCH_RET(schProcessOnTaskFailure(pJob, pTask, code));
14,422,553✔
459
} 
460

461

462
// Note: no more task error processing, handled in function internal
463
int32_t schHandleResponseMsg(SSchJob *pJob, SSchTask *pTask, uint64_t seriesId, int32_t execId, SDataBuf *pMsg, int32_t rspCode) {
1,638,175,882✔
464
  int32_t code = 0;
1,638,175,882✔
465
  int32_t msgType = pMsg->msgType;
1,638,175,882✔
466

467
  bool dropExecNode = (msgType == TDMT_SCH_LINK_BROKEN || SCH_NETWORK_ERR(rspCode));
1,638,189,983✔
468
  if (SCH_IS_QUERY_JOB(pJob)) {
1,638,189,983✔
469
    SCH_ERR_JRET(schUpdateTaskHandle(pJob, pTask, dropExecNode, pMsg->handle, seriesId, execId));
916,557,231✔
470
  }
471
  
472
  SCH_ERR_JRET(schValidateRspMsgType(pJob, pTask, msgType));
1,638,193,243✔
473

474
  if (pTask->seriesId < atomic_load_64(&pJob->seriesId)) {
1,638,174,883✔
475
    SCH_TASK_DLOG("task sId %" PRId64 " is smaller than current job sId %" PRId64, pTask->seriesId, pJob->seriesId);
×
476
    SCH_ERR_JRET(TSDB_CODE_SCH_IGNORE_ERROR);
×
477
  }
478

479
  int32_t reqType = IsReq(pMsg) ? pMsg->msgType : (pMsg->msgType - 1);
1,638,161,375✔
480

481
  if (SCH_DATA_BIND_TASK_NEED_RETRY(pJob, pTask,reqType, rspCode)) {
1,638,165,189✔
482
    SCH_RET(schProcessOnTaskFailure(pJob, pTask, rspCode));
161,190✔
483
  } else if (SCH_JOB_NEED_RETRY(pJob, pTask, reqType, rspCode)) {
1,638,030,900✔
484
    SCH_RET(schHandleJobRetry(pJob, pTask, (SDataBuf *)pMsg, rspCode));
4,517,998✔
485
  }
486

487
  pTask->redirectCtx.inRedirect = false;
1,633,538,052✔
488
  SCH_RET(schProcessResponseMsg(pJob, pTask, pMsg, rspCode));
1,633,516,185✔
489

490
_return:
×
491
  taosMemoryFreeClear(pMsg->pData);
×
492
  SCH_RET(schProcessOnTaskFailure(pJob, pTask, code));
×
493
} 
494

495
int32_t schHandleCallback(void *param, SDataBuf *pMsg, int32_t rspCode) {
1,642,235,929✔
496
  int32_t                code = 0;
1,642,235,929✔
497
  SSchTaskCallbackParam *pParam = (SSchTaskCallbackParam *)param;
1,642,235,929✔
498
  SSchTask              *pTask = NULL;
1,642,235,929✔
499
  SSchJob               *pJob = NULL;
1,642,248,515✔
500

501
  int64_t qid = pParam->queryId;
1,642,251,746✔
502
  qDebug("QID:0x%" PRIx64 ", handle rsp msg, type:%s, handle:%p, code:%s", qid,TMSG_INFO(pMsg->msgType), pMsg->handle,
1,642,250,579✔
503
         tstrerror(rspCode));
504

505
  SCH_ERR_JRET(schProcessOnCbBegin(&pJob, &pTask, pParam->queryId, pParam->refId, pParam->subJobId, pParam->taskId));
1,642,250,703✔
506
  code = schHandleResponseMsg(pJob, pTask, pParam->seriesId, pParam->execId, pMsg, rspCode);
1,638,127,321✔
507
  pMsg->pData = NULL;
1,638,117,785✔
508

509
  schProcessOnCbEnd(pJob, pTask, code);
1,638,131,750✔
510

511
_return:
1,642,260,380✔
512

513
  taosMemoryFreeClear(pMsg->pData);
1,642,271,397✔
514
  taosMemoryFreeClear(pMsg->pEpSet);
1,642,260,816✔
515

516
  qTrace("QID:0x%" PRIx64 ", end to handle rsp msg, type:%s, handle:%p, code:%s", qid, TMSG_INFO(pMsg->msgType), pMsg->handle,
1,642,253,151✔
517
         tstrerror(rspCode));
518

519
  SCH_RET(code);
1,642,237,035✔
520
}
521

522
int32_t schHandleDropCallback(void *param, SDataBuf *pMsg, int32_t code) {
23,796,721✔
523
  SSchTaskCallbackParam *pParam = (SSchTaskCallbackParam *)param;
23,796,721✔
524
  qDebug("QID:0x%" PRIx64 ", SID:0x%" PRIx64 ", CID:0x%" PRIx64 ", TID:0x%" PRIx64 " SJID:%d drop task rsp received, code:0x%x", 
23,796,721✔
525
         pParam->queryId, pParam->seriesId, pParam->clientId, pParam->taskId, pParam->subJobId, code);
526
  // called if drop task rsp received code
527
  if (pMsg->handle == NULL) {
23,797,551✔
528
    qError("sch handle is NULL, may be already released and mem leak");
23,191,240✔
529
  } else {
530
    (void)rpcReleaseHandle(pMsg->handle, TAOS_CONN_CLIENT, 0); // ignore error
606,311✔
531
    pMsg->handle = NULL;
606,311✔
532
  }
533

534
  if (pMsg) {
23,797,633✔
535
    taosMemoryFree(pMsg->pData);
23,797,715✔
536
    taosMemoryFree(pMsg->pEpSet);
23,797,635✔
537
  }
538
  return TSDB_CODE_SUCCESS;
23,797,635✔
539
}
540

541
int32_t schHandleNotifyCallback(void *param, SDataBuf *pMsg, int32_t code) {
2,020✔
542
  SSchTaskCallbackParam *pParam = (SSchTaskCallbackParam *)param;
2,020✔
543
  qDebug("QID:0x%" PRIx64 ", SID:0x%" PRIx64 ", CID:0x%" PRIx64 ", TID:0x%" PRIx64 " SJID:%d task notify rsp received, code:0x%x", 
2,020✔
544
         pParam->queryId, pParam->seriesId, pParam->clientId, pParam->taskId, pParam->subJobId, code);
545
  if (pMsg) {
2,020✔
546
    taosMemoryFreeClear(pMsg->pData);
2,020✔
547
    taosMemoryFreeClear(pMsg->pEpSet);
2,020✔
548
  }
549
  return TSDB_CODE_SUCCESS;
2,020✔
550
}
551

552

553
int32_t schHandleLinkBrokenCallback(void *param, SDataBuf *pMsg, int32_t code) {
4,040✔
554
  SSchCallbackParamHeader *head = (SSchCallbackParamHeader *)param;
4,040✔
555
  (void)rpcReleaseHandle(pMsg->handle, TAOS_CONN_CLIENT, 0); // ignore error
4,040✔
556
  qDebug("handle %p is broken", pMsg->handle);
4,040✔
557
  pMsg->handle = NULL;
4,040✔
558

559
  if (head->isHbParam) {
4,040✔
560
    taosMemoryFreeClear(pMsg->pData);
2,020✔
561
    taosMemoryFreeClear(pMsg->pEpSet);
2,020✔
562

563
    SSchHbCallbackParam *hbParam = (SSchHbCallbackParam *)param;
2,020✔
564
    SSchTrans            trans = {.pTrans = hbParam->pTrans, .pHandle = NULL, .pHandleId = 0};
2,020✔
565
    SCH_ERR_RET(schUpdateHbConnection(&hbParam->nodeEpId, &trans));
2,020✔
566

567
    SCH_ERR_RET(schBuildAndSendHbMsg(&hbParam->nodeEpId, NULL));
2,020✔
568
  } else {
569
    SCH_ERR_RET(schHandleCallback(param, pMsg, code));
2,020✔
570
  }
571

572
  return TSDB_CODE_SUCCESS;
2,020✔
573
}
574

575
int32_t schHandleCommitCallback(void *param, SDataBuf *pMsg, int32_t code) {
4,033,682✔
576
  return schHandleCallback(param, pMsg, code);
4,033,682✔
577
}
578

579
int32_t schHandleHbCallback(void *param, SDataBuf *pMsg, int32_t code) {
284,152,681✔
580
  SSchedulerHbRsp        rsp = {0};
284,152,681✔
581
  SSchHbCallbackParam *pParam = (SSchHbCallbackParam *)param;
284,152,039✔
582

583
  if (code) {
284,152,039✔
584
    qError("hb rsp error:%s", tstrerror(code));
2,425,319✔
585
    (void)rpcReleaseHandle(pMsg->handle, TAOS_CONN_CLIENT, 0); // ignore error
2,425,374✔
586
    pMsg->handle = NULL;
2,425,319✔
587
    SCH_ERR_JRET(code);
2,425,319✔
588
  }
589

590
  if (tDeserializeSSchedulerHbRsp(pMsg->pData, pMsg->len, &rsp)) {
281,726,720✔
591
    qError("invalid hb rsp msg, size:%d", pMsg->len);
2,020✔
592
    SCH_ERR_JRET(TSDB_CODE_QRY_INVALID_INPUT);
2,020✔
593
  }
594

595
  SSchTrans trans = {0};
281,710,930✔
596
  trans.pTrans = pParam->pTrans;
281,712,201✔
597
  trans.pHandle = pMsg->handle;
281,705,978✔
598
  trans.pHandleId = pMsg->handleRefId;
281,720,058✔
599

600
  SCH_ERR_JRET(schUpdateHbConnection(&rsp.epId, &trans));
281,717,900✔
601
  SCH_ERR_JRET(schProcessOnTaskStatusRsp(&rsp.epId, rsp.taskStatus));
281,714,765✔
602

603
_return:
284,146,256✔
604

605
  tFreeSSchedulerHbRsp(&rsp);
284,151,509✔
606
  taosMemoryFree(pMsg->pData);
284,150,785✔
607
  taosMemoryFree(pMsg->pEpSet);
284,151,734✔
608
  SCH_RET(code);
284,148,776✔
609
}
610

611
int32_t schMakeCallbackParam(SSchJob *pJob, SSchTask *pTask, int32_t msgType, bool isHb, SSchTrans *trans,
2,147,483,647✔
612
                             void **pParam) {
613
  SQueryNodeAddr *pAddr = NULL;
2,147,483,647✔
614

615
  if (!isHb) {
2,147,483,647✔
616
    SSchTaskCallbackParam *param = taosMemoryCalloc(1, sizeof(SSchTaskCallbackParam));
2,147,483,647✔
617
    if (NULL == param) {
2,147,483,647✔
UNCOV
618
      SCH_TASK_ELOG("calloc %d failed", (int32_t)sizeof(SSchTaskCallbackParam));
×
UNCOV
619
      SCH_ERR_RET(terrno);
×
620
    }
621

622
    param->queryId = pJob->queryId;
2,147,483,647✔
623
    param->seriesId = pTask->seriesId;
2,147,483,647✔
624
    param->refId = pJob->refId;
2,147,483,647✔
625
    param->clientId = SCH_CLIENT_ID(pTask);
2,147,483,647✔
626
    param->taskId = SCH_TASK_ID(pTask);
2,147,483,647✔
627
    param->subJobId = pJob->subJobId;
2,147,483,647✔
628
    param->pTrans = pJob->conn.pTrans;
2,147,483,647✔
629
    param->execId = pTask->execId;
2,147,483,647✔
630
    *pParam = param;
2,147,483,647✔
631

632
    return TSDB_CODE_SUCCESS;
2,147,483,647✔
633
  }
634

635
  if (TDMT_SCH_LINK_BROKEN == msgType) {
518,448,887✔
636
    SSchHbCallbackParam *param = taosMemoryCalloc(1, sizeof(SSchHbCallbackParam));
260,359,402✔
637
    if (NULL == param) {
258,693,851✔
UNCOV
638
      SCH_TASK_ELOG("calloc %d failed", (int32_t)sizeof(SSchHbCallbackParam));
×
UNCOV
639
      SCH_ERR_RET(terrno);
×
640
    }
641

642
    param->head.isHbParam = true;
258,693,851✔
643

644
    int32_t code = schGetTaskCurrentNodeAddr(pTask, pJob, &pAddr);
259,656,328✔
645
    if (code != TSDB_CODE_SUCCESS) {
260,539,355✔
UNCOV
646
      taosMemoryFree(param);
×
647
      SCH_ERR_RET(code);
12,712✔
648
    }
649

650
    param->nodeEpId.nodeId = pAddr->nodeId;
260,552,067✔
651
    SEp *pEp = SCH_GET_CUR_EP(pAddr);
260,468,291✔
652
    tstrncpy(param->nodeEpId.ep.fqdn, pEp->fqdn, sizeof(param->nodeEpId.ep.fqdn));
259,820,435✔
653
    param->nodeEpId.ep.port = pEp->port;
260,245,499✔
654
    param->pTrans = trans->pTrans;
259,707,580✔
655
    *pParam = param;
259,687,974✔
656

657
    return TSDB_CODE_SUCCESS;
259,656,588✔
658
  }
659

660
  // hb msg
661
  SSchHbCallbackParam *param = taosMemoryCalloc(1, sizeof(SSchHbCallbackParam));
258,089,485✔
662
  if (NULL == param) {
256,837,529✔
UNCOV
663
    qError("calloc SSchTaskCallbackParam failed");
×
UNCOV
664
    SCH_ERR_RET(terrno);
×
665
  }
666

667
  param->head.isHbParam = true;
256,837,529✔
668
  param->pTrans = trans->pTrans;
257,080,328✔
669
  *pParam = param;
257,445,006✔
670

671
  return TSDB_CODE_SUCCESS;
257,570,329✔
672
}
673

674
int32_t schGenerateCallBackInfo(SSchJob *pJob, SSchTask *pTask, void *msg, uint32_t msgSize, int32_t msgType,
2,147,483,647✔
675
                                SSchTrans *trans, bool isHb, SMsgSendInfo **pMsgSendInfo) {
676
  int32_t       code = 0;
2,147,483,647✔
677
  SMsgSendInfo *msgSendInfo = taosMemoryCalloc(1, sizeof(SMsgSendInfo));
2,147,483,647✔
678
  if (NULL == msgSendInfo) {
2,147,483,647✔
UNCOV
679
    qError("calloc SMsgSendInfo size %d failed", (int32_t)sizeof(SMsgSendInfo));
×
UNCOV
680
    SCH_ERR_JRET(terrno);
×
681
  }
682

683
  msgSendInfo->paramFreeFp = taosAutoMemoryFree;
2,147,483,647✔
684
  SCH_ERR_JRET(schMakeCallbackParam(pJob, pTask, msgType, isHb, trans, &msgSendInfo->param));
2,147,483,647✔
685

686
  SCH_ERR_JRET(schGetCallbackFp(msgType, &msgSendInfo->fp));
2,147,483,647✔
687

688
  if (pJob) {
2,147,483,647✔
689
    msgSendInfo->requestId = pJob->conn.requestId;
2,147,483,647✔
690
    msgSendInfo->requestObjRefId = pJob->conn.requestObjRefId;
2,147,483,647✔
691
  } else {
692
    SCH_ERR_JRET(taosGetSystemUUIDU64(&msgSendInfo->requestId));
257,689,744✔
693
  }
694

695
  qDebug("ahandle %p alloced, QID:0x%" PRIx64, msgSendInfo, msgSendInfo->requestId);
2,147,483,647✔
696

697
  if (TDMT_SCH_LINK_BROKEN != msgType) {
2,147,483,647✔
698
    msgSendInfo->msgInfo.pData = msg;
2,147,483,647✔
699
    msgSendInfo->msgInfo.len = msgSize;
2,147,483,647✔
700
    msgSendInfo->msgInfo.handle = trans->pHandle;
2,147,483,647✔
701
    msgSendInfo->msgType = msgType;
2,147,483,647✔
702
  }
703

704
  *pMsgSendInfo = msgSendInfo;
2,147,483,647✔
705

706
  return TSDB_CODE_SUCCESS;
2,147,483,647✔
707

UNCOV
708
_return:
×
709

UNCOV
710
  if (msgSendInfo) {
×
711
    destroySendMsgInfo(msgSendInfo);
×
712
  }
713

714
  taosMemoryFree(msg);
×
715

UNCOV
716
  SCH_RET(code);
×
717
}
718

719
int32_t schGetCallbackFp(int32_t msgType, __async_send_cb_fn_t *fp) {
2,147,483,647✔
720
  switch (msgType) {
2,147,483,647✔
721
    case TDMT_VND_CREATE_TABLE:
2,105,330,390✔
722
    case TDMT_VND_DROP_TABLE:
723
    case TDMT_VND_ALTER_TABLE:
724
    case TDMT_VND_SUBMIT:
725
    case TDMT_SCH_QUERY:
726
    case TDMT_SCH_MERGE_QUERY:
727
    case TDMT_VND_DELETE:
728
    case TDMT_SCH_EXPLAIN:
729
    case TDMT_SCH_FETCH:
730
    case TDMT_SCH_MERGE_FETCH:
731
      *fp = schHandleCallback;
2,105,330,390✔
732
      break;
2,105,387,020✔
733
    case TDMT_SCH_DROP_TASK:
552,059,860✔
734
      *fp = schHandleDropCallback;
552,059,860✔
735
      break;
552,060,869✔
736
    case TDMT_SCH_TASK_NOTIFY:
33,938✔
737
      *fp = schHandleNotifyCallback;
33,938✔
738
      break;
33,938✔
739
    case TDMT_SCH_QUERY_HEARTBEAT:
517,777,914✔
740
      *fp = schHandleHbCallback;
517,777,914✔
741
      break;
518,173,700✔
742
    case TDMT_VND_COMMIT:
4,032,353✔
743
      *fp = schHandleCommitCallback;
4,032,353✔
744
      break;
4,033,569✔
745
    case TDMT_SCH_LINK_BROKEN:
813,953,992✔
746
      *fp = schHandleLinkBrokenCallback;
813,953,992✔
747
      break;
813,882,703✔
748
    default:
110,858✔
749
      qError("unknown msg type for callback, msgType:%d", msgType);
110,858✔
750
      SCH_ERR_RET(TSDB_CODE_APP_ERROR);
155✔
751
  }
752

753
  return TSDB_CODE_SUCCESS;
2,147,483,647✔
754
}
755

756
/*
757
int32_t schMakeHbCallbackParam(SSchJob *pJob, SSchTask *pTask, void **pParam) {
758
  SSchHbCallbackParam *param = taosMemoryCalloc(1, sizeof(SSchHbCallbackParam));
759
  if (NULL == param) {
760
    SCH_TASK_ELOG("calloc %d failed", (int32_t)sizeof(SSchHbCallbackParam));
761
    SCH_ERR_RET(TSDB_CODE_OUT_OF_MEMORY);
762
  }
763

764
  param->head.isHbParam = true;
765

766
  SQueryNodeAddr *addr = taosArrayGet(pTask->candidateAddrs, pTask->candidateIdx);
767

768
  param->nodeEpId.nodeId = addr->nodeId;
769
  SEp* pEp = SCH_GET_CUR_EP(addr);
770
  tstrncpy(param->nodeEpId.ep.fqdn, pEp->fqdn, sizeof(param->nodeEpId.ep.fqdn));
771
  param->nodeEpId.ep.port = pEp->port;
772
  param->pTrans = pJob->pTrans;
773

774
  *pParam = param;
775

776
  return TSDB_CODE_SUCCESS;
777
}
778
*/
779

780
int32_t schCloneHbRpcCtx(SRpcCtx *pSrc, SRpcCtx *pDst) {
258,174,889✔
781
  int32_t code = 0;
258,174,889✔
782
  TAOS_MEMCPY(pDst, pSrc, sizeof(SRpcCtx));
258,174,889✔
783
  pDst->brokenVal.val = NULL;
258,174,889✔
784
  pDst->args = NULL;
258,068,800✔
785

786
  SCH_ERR_RET(schCloneSMsgSendInfo(pSrc->brokenVal.val, &pDst->brokenVal.val));
257,543,127✔
787

788
  pDst->args = taosHashInit(1, taosGetDefaultHashFunction(TSDB_DATA_TYPE_INT), false, HASH_ENTRY_LOCK);
257,172,939✔
789
  if (NULL == pDst->args) {
258,401,832✔
UNCOV
790
    qError("taosHashInit %d RpcCtx failed", 1);
×
UNCOV
791
    SCH_ERR_JRET(terrno);
×
792
  }
793

794
  SRpcCtxVal dst = {0};
258,462,917✔
795
  void      *pIter = taosHashIterate(pSrc->args, NULL);
258,262,234✔
796
  while (pIter) {
516,788,789✔
797
    SRpcCtxVal *pVal = (SRpcCtxVal *)pIter;
258,155,027✔
798
    int32_t    *msgType = taosHashGetKey(pIter, NULL);
258,155,027✔
799

800
    dst = *pVal;
258,270,102✔
801
    dst.val = NULL;
258,355,274✔
802

803
    SCH_ERR_JRET(schCloneSMsgSendInfo(pVal->val, &dst.val));
258,355,274✔
804

805
    if (taosHashPut(pDst->args, msgType, sizeof(*msgType), &dst, sizeof(dst))) {
258,299,893✔
UNCOV
806
      qError("taosHashPut msg %d to rpcCtx failed", *msgType);
×
UNCOV
807
      (*pSrc->freeFunc)(dst.val);
×
UNCOV
808
      SCH_ERR_JRET(TSDB_CODE_OUT_OF_MEMORY);
×
809
    }
810

811
    pIter = taosHashIterate(pSrc->args, pIter);
258,650,784✔
812
  }
813

814
  return TSDB_CODE_SUCCESS;
258,633,762✔
815

UNCOV
816
_return:
×
817

UNCOV
818
  schFreeRpcCtx(pDst);
×
819
  SCH_RET(code);
×
820
}
821

822
int32_t schMakeHbRpcCtx(SSchJob *pJob, SSchTask *pTask, SRpcCtx *pCtx) {
260,424,459✔
823
  int32_t              code = 0;
260,424,459✔
824
  SSchHbCallbackParam *param = NULL;
260,424,459✔
825
  SMsgSendInfo        *pMsgSendInfo = NULL;
260,424,459✔
826
  SQueryNodeAddr      *pAddr = NULL;
260,424,459✔
827
  SQueryNodeEpId       epId = {0};
260,389,692✔
828
  int32_t              msgType = TDMT_SCH_QUERY_HEARTBEAT_RSP;
260,666,037✔
829

830
  code = schGetTaskCurrentNodeAddr(pTask, pJob, &pAddr);
260,504,184✔
831
  if (code != TSDB_CODE_SUCCESS) {
260,639,375✔
UNCOV
832
    SCH_ERR_JRET(code);
×
833
  }
834

835
  epId.nodeId = pAddr->nodeId;
260,639,375✔
836
  TAOS_MEMCPY(&epId.ep, SCH_GET_CUR_EP(pAddr), sizeof(SEp));
260,409,615✔
837

838
  pCtx->args = taosHashInit(1, taosGetDefaultHashFunction(TSDB_DATA_TYPE_INT), false, HASH_ENTRY_LOCK);
260,036,230✔
839
  if (NULL == pCtx->args) {
260,670,539✔
UNCOV
840
    SCH_TASK_ELOG("taosHashInit %d RpcCtx failed", 1);
×
UNCOV
841
    SCH_ERR_RET(terrno);
×
842
  }
843

844
  pMsgSendInfo = taosMemoryCalloc(1, sizeof(SMsgSendInfo));
260,726,721✔
845
  if (NULL == pMsgSendInfo) {
259,669,857✔
UNCOV
846
    SCH_TASK_ELOG("calloc %d failed", (int32_t)sizeof(SMsgSendInfo));
×
UNCOV
847
    SCH_ERR_JRET(terrno);
×
848
  }
849

850
  param = taosMemoryCalloc(1, sizeof(SSchHbCallbackParam));
259,669,857✔
851
  if (NULL == param) {
258,929,753✔
UNCOV
852
    SCH_TASK_ELOG("calloc %d failed", (int32_t)sizeof(SSchHbCallbackParam));
×
UNCOV
853
    SCH_ERR_JRET(terrno);
×
854
  }
855

856
  __async_send_cb_fn_t fp = NULL;
258,929,753✔
857
  SCH_ERR_JRET(schGetCallbackFp(TDMT_SCH_QUERY_HEARTBEAT, &fp));
259,535,644✔
858

859
  param->head.isHbParam = true;
260,039,538✔
860
  param->nodeEpId = epId;
260,148,968✔
861
  param->pTrans = pJob->conn.pTrans;
260,549,833✔
862

863
  pMsgSendInfo->param = param;
260,107,290✔
864
  pMsgSendInfo->paramFreeFp = taosAutoMemoryFree;
259,684,943✔
865
  pMsgSendInfo->fp = fp;
260,654,979✔
866

867
  SRpcCtxVal ctxVal = {.val = pMsgSendInfo, .clone = schCloneSMsgSendInfo};
260,071,415✔
868
  if (taosHashPut(pCtx->args, &msgType, sizeof(msgType), &ctxVal, sizeof(ctxVal))) {
260,158,875✔
UNCOV
869
    SCH_TASK_ELOG("taosHashPut msg %d to rpcCtx failed", msgType);
×
UNCOV
870
    SCH_ERR_JRET(TSDB_CODE_OUT_OF_MEMORY);
×
871
  }
872

873
  SCH_ERR_JRET(schMakeBrokenLinkVal(pJob, pTask, &pCtx->brokenVal, true));
261,065,882✔
874
  pCtx->freeFunc = schFreeRpcCtxVal;
261,028,589✔
875

876
  return TSDB_CODE_SUCCESS;
260,975,992✔
877

UNCOV
878
_return:
×
879

UNCOV
880
  taosHashCleanup(pCtx->args);
×
881
  taosMemoryFreeClear(param);
×
UNCOV
882
  taosMemoryFreeClear(pMsgSendInfo);
×
883

884
  SCH_RET(code);
×
885
}
886

887
int32_t schMakeBrokenLinkVal(SSchJob *pJob, SSchTask *pTask, SRpcBrokenlinkVal *brokenVal, bool isHb) {
814,060,089✔
888
  int32_t       code = 0;
814,060,089✔
889
  int32_t       msgType = TDMT_SCH_LINK_BROKEN;
814,060,089✔
890
  SSchTrans     trans = {.pTrans = pJob->conn.pTrans};
814,060,089✔
891
  SMsgSendInfo *pMsgSendInfo = NULL;
814,137,314✔
892
  SCH_ERR_JRET(schGenerateCallBackInfo(pJob, pTask, NULL, 0, msgType, &trans, isHb, &pMsgSendInfo));
813,652,185✔
893

894
  brokenVal->msgType = msgType;
814,847,278✔
895
  brokenVal->val = pMsgSendInfo;
814,866,740✔
896
  brokenVal->clone = schCloneSMsgSendInfo;
814,871,600✔
897

898
  return TSDB_CODE_SUCCESS;
814,840,617✔
899

UNCOV
900
_return:
×
901

UNCOV
902
  taosMemoryFreeClear(pMsgSendInfo->param);
×
903
  taosMemoryFreeClear(pMsgSendInfo);
×
904

905
  SCH_RET(code);
×
906
}
907

908
int32_t schMakeQueryRpcCtx(SSchJob *pJob, SSchTask *pTask, SRpcCtx *pCtx) {
553,468,718✔
909
  int32_t       code = 0;
553,468,718✔
910
  SMsgSendInfo *pExplainMsgSendInfo = NULL;
553,468,718✔
911

912
  pCtx->args = taosHashInit(1, taosGetDefaultHashFunction(TSDB_DATA_TYPE_INT), false, HASH_ENTRY_LOCK);
553,513,731✔
913
  if (NULL == pCtx->args) {
553,571,452✔
UNCOV
914
    SCH_TASK_ELOG("taosHashInit %d RpcCtx failed", 1);
×
UNCOV
915
    SCH_ERR_RET(terrno);
×
916
  }
917

918
  SSchTrans trans = {.pTrans = pJob->conn.pTrans, .pHandle = SCH_GET_TASK_HANDLE(pTask)};
553,574,591✔
919
  SCH_ERR_JRET(schGenerateCallBackInfo(pJob, pTask, NULL, 0, TDMT_SCH_EXPLAIN, &trans, false, &pExplainMsgSendInfo));
553,649,767✔
920

921
  int32_t    msgType = TDMT_SCH_EXPLAIN_RSP;
553,748,441✔
922
  SRpcCtxVal ctxVal = {.val = pExplainMsgSendInfo, .clone = schCloneSMsgSendInfo};
553,730,267✔
923
  if (taosHashPut(pCtx->args, &msgType, sizeof(msgType), &ctxVal, sizeof(ctxVal))) {
553,791,048✔
UNCOV
924
    SCH_TASK_ELOG("taosHashPut msg %d to rpcCtx failed", msgType);
×
UNCOV
925
    SCH_ERR_JRET(TSDB_CODE_OUT_OF_MEMORY);
×
926
  }
927

928
  SCH_ERR_JRET(schMakeBrokenLinkVal(pJob, pTask, &pCtx->brokenVal, false));
553,861,148✔
929
  pCtx->freeFunc = schFreeRpcCtxVal;
553,796,905✔
930

931
  return TSDB_CODE_SUCCESS;
553,774,904✔
932

UNCOV
933
_return:
×
934

UNCOV
935
  taosHashCleanup(pCtx->args);
×
936

UNCOV
937
  if (pExplainMsgSendInfo) {
×
938
    taosMemoryFreeClear(pExplainMsgSendInfo->param);
×
UNCOV
939
    taosMemoryFreeClear(pExplainMsgSendInfo);
×
940
  }
941

942
  SCH_RET(code);
×
943
}
944

945
int32_t schCloneCallbackParam(SSchCallbackParamHeader *pSrc, SSchCallbackParamHeader **pDst) {
628,391,820✔
946
  if (pSrc->isHbParam) {
628,391,820✔
947
    SSchHbCallbackParam *dst = taosMemoryMalloc(sizeof(SSchHbCallbackParam));
542,228,742✔
948
    if (NULL == dst) {
539,878,943✔
UNCOV
949
      qError("malloc SSchHbCallbackParam failed");
×
UNCOV
950
      SCH_ERR_RET(terrno);
×
951
    }
952

953
    TAOS_MEMCPY(dst, pSrc, sizeof(*dst));
539,878,943✔
954
    *pDst = (SSchCallbackParamHeader *)dst;
539,878,943✔
955

956
    return TSDB_CODE_SUCCESS;
542,217,644✔
957
  }
958

959
  SSchTaskCallbackParam *dst = taosMemoryMalloc(sizeof(SSchTaskCallbackParam));
86,292,466✔
960
  if (NULL == dst) {
86,279,002✔
UNCOV
961
    qError("malloc SSchTaskCallbackParam failed");
×
UNCOV
962
    SCH_ERR_RET(terrno);
×
963
  }
964

965
  TAOS_MEMCPY(dst, pSrc, sizeof(*dst));
86,279,002✔
966
  *pDst = (SSchCallbackParamHeader *)dst;
86,279,002✔
967

968
  return TSDB_CODE_SUCCESS;
86,278,850✔
969
}
970

971
int32_t schCloneSMsgSendInfo(void *src, void **dst) {
628,034,864✔
972
  SMsgSendInfo *pSrc = src;
628,034,864✔
973
  int32_t       code = 0;
628,034,864✔
974
  SMsgSendInfo *pDst = taosMemoryCalloc(1, sizeof(*pSrc));
628,034,864✔
975
  if (NULL == pDst) {
626,503,053✔
UNCOV
976
    qError("malloc SMsgSendInfo for rpcCtx failed, len:%d", (int32_t)sizeof(*pSrc));
×
UNCOV
977
    SCH_ERR_RET(terrno);
×
978
  }
979

980
  TAOS_MEMCPY(pDst, pSrc, sizeof(*pSrc));
626,503,053✔
981
  pDst->param = NULL;
626,503,053✔
982

983
  SCH_ERR_JRET(schCloneCallbackParam(pSrc->param, (SSchCallbackParamHeader **)&pDst->param));
628,520,760✔
984
  pDst->paramFreeFp = taosAutoMemoryFree;
628,353,249✔
985

986
  *dst = pDst;
628,426,179✔
987

988
  return TSDB_CODE_SUCCESS;
628,486,876✔
989

UNCOV
990
_return:
×
991

UNCOV
992
  taosMemoryFreeClear(pDst);
×
993
  SCH_RET(code);
×
994
}
995

996
int32_t schUpdateSendTargetInfo(SMsgSendInfo *pMsgSendInfo, SQueryNodeAddr *addr, SSchTask *pTask) {
2,147,483,647✔
997
  if (NULL == pTask || addr->nodeId < MNODE_HANDLE) {
2,147,483,647✔
998
    return TSDB_CODE_SUCCESS;
261,762,846✔
999
  }
1000

1001
  if (addr->nodeId == MNODE_HANDLE) {
2,104,676,939✔
1002
    pMsgSendInfo->target.type = TARGET_TYPE_MNODE;
21,211,958✔
1003
  } else {
1004
    pMsgSendInfo->target.type = TARGET_TYPE_VNODE;
2,083,486,406✔
1005
    pMsgSendInfo->target.vgId = addr->nodeId;
2,083,518,175✔
1006
    pMsgSendInfo->target.dbFName = taosStrdup(pTask->plan->dbFName);
2,083,515,867✔
1007
    if (NULL == pMsgSendInfo->target.dbFName) {
2,083,328,046✔
UNCOV
1008
      return terrno;
×
1009
    }
1010
  }
1011

1012
  return TSDB_CODE_SUCCESS;
2,104,449,872✔
1013
}
1014

1015
int32_t schAsyncSendMsg(SSchJob *pJob, SSchTask *pTask, SSchTrans *trans, SQueryNodeAddr *addr, int32_t msgType,
2,147,483,647✔
1016
                        void *msg, uint32_t msgSize, bool persistHandle, SRpcCtx *ctx) {
1017
  int32_t code = 0;
2,147,483,647✔
1018
  SEpSet *epSet = &addr->epSet;
2,147,483,647✔
1019

1020
  SMsgSendInfo *pMsgSendInfo = NULL;
2,147,483,647✔
1021
  bool          isHb = (TDMT_SCH_QUERY_HEARTBEAT == msgType);
2,147,483,647✔
1022
  SCH_ERR_JRET(schGenerateCallBackInfo(pJob, pTask, msg, msgSize, msgType, trans, isHb, &pMsgSendInfo));
2,147,483,647✔
1023
  SCH_ERR_JRET(schUpdateSendTargetInfo(pMsgSendInfo, addr, pTask));
2,147,483,647✔
1024

1025
  if (isHb && persistHandle && trans->pHandle == 0) {
2,147,483,647✔
1026
    int64_t refId = 0;
258,599,661✔
1027
    code = rpcAllocHandle(&refId); 
258,599,396✔
1028
    if (code != 0) {
258,562,770✔
UNCOV
1029
      SCH_TASK_ELOG("rpcAllocHandle failed, code:%x", code);
×
1030
      SCH_ERR_JRET(code);
5,170✔
1031
    }
1032
    trans->pHandle = (void *)refId;
258,567,940✔
1033
    pMsgSendInfo->msgInfo.handle =trans->pHandle;
258,575,224✔
1034
  } 
1035

1036
  if (pJob && pTask) {
2,147,483,647✔
1037
    SCH_TASK_DLOG("start to send %s msg to node[%d,%s,%d], pTrans:%p, pHandle:%p", TMSG_INFO(msgType), addr->nodeId,
2,107,741,515✔
1038
           epSet->eps[epSet->inUse].fqdn, epSet->eps[epSet->inUse].port, trans->pTrans, trans->pHandle);
1039
  } else {
1040
    qDebug("start to send %s msg to node[%d,%s,%d], pTrans:%p, pHandle:%p", TMSG_INFO(msgType), addr->nodeId,
258,505,968✔
1041
           epSet->eps[epSet->inUse].fqdn, epSet->eps[epSet->inUse].port, trans->pTrans, trans->pHandle);
1042
  }
1043
  
1044
  if (pTask) {
2,147,483,647✔
1045
    pTask->lastMsgType = msgType;
2,107,622,122✔
1046
  }
1047

1048
  code = asyncSendMsgToServerExt(trans->pTrans, epSet, NULL, pMsgSendInfo, persistHandle, ctx);
2,147,483,647✔
1049
  pMsgSendInfo = NULL;
2,147,483,647✔
1050
  if (code) {
2,147,483,647✔
1051
    SCH_ERR_JRET(code);
367,927✔
1052
  }
1053

1054
  if (pJob) {
2,147,483,647✔
1055
    SCH_TASK_TLOG("req msg sent, type:%d, %s", msgType, TMSG_INFO(msgType));
2,107,696,962✔
1056
  } else {
1057
    qTrace("req msg sent, type:%d, %s", msgType, TMSG_INFO(msgType));
258,540,901✔
1058
  }
1059
  return TSDB_CODE_SUCCESS;
2,147,483,647✔
1060

1061
_return:
399,853✔
1062

1063
  if (pJob) {
367,927✔
1064
    SCH_TASK_ELOG("fail to send msg, type:%d, %s, error:%s", msgType, TMSG_INFO(msgType), tstrerror(code));
363,609✔
1065
  } else {
1066
    qError("fail to send msg, type:%d, %s, error:%s", msgType, TMSG_INFO(msgType), tstrerror(code));
4,318✔
1067
  }
1068

1069
  if (pMsgSendInfo) {
367,927✔
UNCOV
1070
    destroySendMsgInfo(pMsgSendInfo);
×
1071
  }
1072

1073
  SCH_RET(code);
367,927✔
1074
}
1075

1076
int32_t schBuildAndSendHbMsg(SQueryNodeEpId *nodeEpId, SArray *taskAction) {
258,756,194✔
1077
  SSchedulerHbReq req = {0};
258,756,194✔
1078
  int32_t         code = 0;
258,756,194✔
1079
  SRpcCtx         rpcCtx = {0};
258,756,194✔
1080
  SSchTrans       trans = {0};
258,756,194✔
1081
  int32_t         msgType = TDMT_SCH_QUERY_HEARTBEAT;
258,756,194✔
1082

1083
  req.header.vgId = nodeEpId->nodeId;
258,756,194✔
1084
  req.clientId = schMgmt.clientId;
258,756,194✔
1085
  TAOS_MEMCPY(&req.epId, nodeEpId, sizeof(SQueryNodeEpId));
258,756,194✔
1086

1087
  SCH_LOCK(SCH_READ, &schMgmt.hbLock);
258,756,194✔
1088
  SSchHbTrans *hb = taosHashGet(schMgmt.hbConnections, nodeEpId, sizeof(SQueryNodeEpId));
258,691,362✔
1089
  if (NULL == hb) {
258,295,443✔
1090
    SCH_UNLOCK(SCH_READ, &schMgmt.hbLock);
90,479✔
1091
    qError("hb connection no longer exist, nodeId:%d, fqdn:%s, port:%d", nodeEpId->nodeId, nodeEpId->ep.fqdn,
90,561✔
1092
           nodeEpId->ep.port);
1093
    return TSDB_CODE_SUCCESS;
90,561✔
1094
  }
1095

1096
  SCH_LOCK(SCH_WRITE, &hb->lock);
258,204,964✔
1097
  code = schCloneHbRpcCtx(&hb->rpcCtx, &rpcCtx);
257,775,621✔
1098
  TAOS_MEMCPY(&trans, &hb->trans, sizeof(trans));
258,340,471✔
1099
  if (NULL == hb->trans.pTrans) {
258,296,688✔
UNCOV
1100
    qError("NULL pTrans got from hbConnections for epId:%d", nodeEpId->nodeId);
×
1101
  }
1102
  SCH_UNLOCK(SCH_WRITE, &hb->lock);
258,468,035✔
1103
  SCH_UNLOCK(SCH_READ, &schMgmt.hbLock);
258,153,635✔
1104

1105
  SCH_ERR_RET(code);
257,082,109✔
1106

1107
  int32_t msgSize = tSerializeSSchedulerHbReq(NULL, 0, &req);
257,082,109✔
1108
  if (msgSize < 0) {
257,963,233✔
UNCOV
1109
    qError("tSerializeSSchedulerHbReq hbReq failed, size:%d", msgSize);
×
UNCOV
1110
    SCH_ERR_JRET(TSDB_CODE_OUT_OF_MEMORY);
×
1111
  }
1112
  void *msg = taosMemoryCalloc(1, msgSize);
257,963,233✔
1113
  if (NULL == msg) {
258,050,786✔
UNCOV
1114
    qError("calloc hb req %d failed", msgSize);
×
UNCOV
1115
    SCH_ERR_JRET(terrno);
×
1116
  }
1117

1118
  if (tSerializeSSchedulerHbReq(msg, msgSize, &req) < 0) {
258,050,786✔
UNCOV
1119
    qError("tSerializeSSchedulerHbReq hbReq failed, size:%d", msgSize);
×
UNCOV
1120
    SCH_ERR_JRET(TSDB_CODE_OUT_OF_MEMORY);
×
1121
  }
1122

1123
  int64_t        transporterId = 0;
258,018,701✔
1124
  SQueryNodeAddr addr = {.nodeId = nodeEpId->nodeId};
258,018,701✔
1125
  addr.epSet.inUse = 0;
258,150,035✔
1126
  addr.epSet.numOfEps = 1;
258,150,035✔
1127
  TAOS_MEMCPY(&addr.epSet.eps[0], &nodeEpId->ep, sizeof(nodeEpId->ep));
258,150,035✔
1128

1129
  code = schAsyncSendMsg(NULL, NULL, &trans, &addr, msgType, msg, msgSize, true, &rpcCtx);
257,093,860✔
1130
  msg = NULL;
258,618,443✔
1131
  SCH_ERR_JRET(code);
258,618,443✔
1132

1133
  return TSDB_CODE_SUCCESS;
258,614,125✔
1134

1135
_return:
4,318✔
1136

1137
  taosMemoryFreeClear(msg);
4,318✔
1138
  schFreeRpcCtx(&rpcCtx);
4,318✔
1139
  SCH_RET(code);
4,318✔
1140
}
1141

1142
int32_t schBuildSubJobEndpoints(SSchTask *pTask, SArray** ppRes, SSchJob* pTarget) {
553,650,109✔
1143
  *ppRes = NULL;
553,650,109✔
1144
  int32_t code = TSDB_CODE_SUCCESS;
553,655,818✔
1145
  SSchJob* pJob = NULL;
553,655,818✔
1146
  
1147
  if (SCH_IS_PARENT_JOB(pTarget) && !SCH_JOB_GOT_SUB_JOBS(pTarget)) {
553,655,818✔
1148
    pJob = pTarget;
333,385,167✔
1149
    SCH_TASK_DLOG("task query without subEndPoints, pJob:%p", pTarget);
333,385,167✔
1150
    return TSDB_CODE_SUCCESS;
333,411,363✔
1151
  }
1152
  
1153
  int32_t subJobNum = SCH_IS_PARENT_JOB(pTarget) ? pTarget->subJobs->size : pTarget->subJobId;
220,293,023✔
1154
  SSchJob* pParent = SCH_IS_PARENT_JOB(pTarget) ? pTarget : (SSchJob*)pTarget->parent;
220,232,628✔
1155
  SDownstreamSourceNode* pSource = NULL;
220,227,455✔
1156

1157
  if (subJobNum > 0) {
220,210,858✔
1158
    *ppRes = taosArrayInit(subJobNum, POINTER_BYTES);
133,108,126✔
1159
    if (NULL == *ppRes) {
133,115,435✔
UNCOV
1160
      SCH_ERR_RET(terrno);
×
1161
    }
1162
  }
1163
  
1164
  for (int32_t i = 0; i < subJobNum; ++i) {
462,651,701✔
1165
    pJob = taosArrayGetP(pParent->subJobs, i);
242,334,875✔
1166
    if (NULL == pJob || NULL == pJob->fetchTask) {
242,345,559✔
UNCOV
1167
      SCH_JOB_ELOG("no fetchTask set in job, pJob:%p, fetchTask:%p", pJob, pJob ? pJob->fetchTask : NULL);
×
UNCOV
1168
      SCH_ERR_JRET(TSDB_CODE_SCH_INTERNAL_ERROR);
×
1169
    }
1170

1171
    SCH_ERR_JRET(nodesMakeNode(QUERY_NODE_DOWNSTREAM_SOURCE, (SNode**)&pSource));
242,351,278✔
1172
    
1173
    memcpy(&pSource->addr, &pJob->resNode, sizeof(pSource->addr));
242,308,468✔
1174
    pSource->clientId = pJob->fetchTask->clientId;
242,303,323✔
1175
    pSource->taskId = pJob->fetchTask->taskId;
242,301,587✔
1176
    pSource->sId = pJob->fetchTask->seriesId;
242,296,442✔
1177
    pSource->execId = pJob->fetchTask->execId;
242,303,316✔
1178
    pSource->fetchMsgType = SCH_FETCH_TYPE(pJob->fetchTask);
242,306,781✔
1179
    pSource->localExec = pJob->attr.localExec;
242,312,066✔
1180
    if (NULL == taosArrayPush(*ppRes, &pSource)) {
484,653,391✔
UNCOV
1181
      nodesDestroyNode((SNode *)pSource);
×
UNCOV
1182
      SCH_ERR_JRET(terrno);
×
1183
    }
1184
  }
1185

1186
  pJob = pTarget;
220,316,826✔
1187
  SCH_TASK_DLOG("task query with %d subEndPoints", subJobNum);
220,316,826✔
1188

1189
_return:
220,316,826✔
1190

1191
  if (code) {
220,286,650✔
UNCOV
1192
    taosArrayDestroyP(*ppRes, (FDelete)nodesDestroyNode);
×
UNCOV
1193
    *ppRes = NULL;
×
1194
  }
1195
  
1196
  return code;
220,264,334✔
1197
}
1198

1199
int32_t schBuildAndSendMsg(SSchJob *pJob, SSchTask *pTask, SQueryNodeAddr *addr, int32_t msgType, void* param) {
2,107,640,172✔
1200
  int32_t  msgSize = 0;
2,107,640,172✔
1201
  void    *msg = NULL;
2,107,640,172✔
1202
  int32_t  code = 0;
2,107,640,172✔
1203
  bool     isCandidateAddr = false;
2,107,640,172✔
1204
  bool     persistHandle = false;
2,107,640,172✔
1205
  SRpcCtx  rpcCtx = {0};
2,107,640,172✔
1206

1207
  if (NULL == addr) {
2,107,657,222✔
1208
    code = schGetTaskCurrentNodeAddr(pTask, pJob, &addr);
1,274,956,521✔
1209
    if (code != TSDB_CODE_SUCCESS) {
1,275,181,390✔
1210
      SCH_ERR_JRET(code);
2,020✔
1211
    }
1212
    
1213
    isCandidateAddr = true;
1,275,179,370✔
1214
    SCH_TASK_TLOG("target candidateIdx %d, epInUse %d/%d", pTask->candidateIdx, addr->epSet.inUse,
1,275,179,370✔
1215
                  addr->epSet.numOfEps);
1216
  }
1217

1218
  switch (msgType) {
2,107,369,966✔
1219
    case TDMT_VND_CREATE_TABLE:
719,500,749✔
1220
    case TDMT_VND_DROP_TABLE:
1221
    case TDMT_VND_ALTER_TABLE:
1222
    case TDMT_VND_SUBMIT:
1223
    case TDMT_VND_COMMIT: {
1224
      msgSize = pTask->msgLen;
719,500,749✔
1225
      msg = pTask->msg;
719,574,358✔
1226
      pTask->msg = NULL;
719,548,784✔
1227
      break;
719,621,107✔
1228
    }
1229

1230
    case TDMT_VND_DELETE: {
1,925,712✔
1231
      SVDeleteReq req = {0};
1,925,712✔
1232
      req.header.vgId = addr->nodeId;
1,925,712✔
1233
      req.sId = pTask->seriesId;
1,925,712✔
1234
      req.queryId = pJob->queryId;
1,925,712✔
1235
      req.clientId = pTask->clientId;
1,925,712✔
1236
      req.taskId = pTask->taskId;
1,925,712✔
1237
      req.phyLen = pTask->msgLen;
1,925,712✔
1238
      req.sqlLen = strlen(pJob->sql);
1,925,712✔
1239
      req.sql = (char *)pJob->sql;
1,925,712✔
1240
      req.msg = pTask->msg;
1,925,712✔
1241
      req.source       = pJob->source;
1,925,712✔
1242
      req.secureDelete = pJob->secureDelete;
1,925,712✔
1243
      msgSize = tSerializeSVDeleteReq(NULL, 0, &req);
1,925,712✔
1244
      if (msgSize < 0) {
1,925,712✔
UNCOV
1245
        SCH_TASK_ELOG("tSerializeSVDeleteReq failed, code:%x", terrno);
×
UNCOV
1246
        SCH_ERR_JRET(terrno);
×
1247
      }
1248
      msg = taosMemoryCalloc(1, msgSize);
1,925,712✔
1249
      if (NULL == msg) {
1,925,712✔
UNCOV
1250
        SCH_TASK_ELOG("calloc %d failed", msgSize);
×
UNCOV
1251
        SCH_ERR_JRET(terrno);
×
1252
      }
1253

1254
      msgSize = tSerializeSVDeleteReq(msg, msgSize, &req);
1,925,712✔
1255
      if (msgSize < 0) {
1,925,712✔
UNCOV
1256
        SCH_TASK_ELOG("tSerializeSVDeleteReq second failed, code:%x", terrno);
×
UNCOV
1257
        SCH_ERR_JRET(terrno);
×
1258
      }
1259
      break;
1,925,712✔
1260
    }
1261
    case TDMT_SCH_QUERY:
553,258,725✔
1262
    case TDMT_SCH_MERGE_QUERY: {
1263
      int32_t newPhase = (TDMT_SCH_QUERY == msgType) ? QUERY_PHASE_EXEC_DATA_QUERY : QUERY_PHASE_EXEC_MERGE_QUERY;
553,258,725✔
1264
      SCH_UPDATE_JOB_PHASE_IF_CHANGED(pJob, newPhase);
948,372,994✔
1265

1266
      SCH_ERR_RET(schMakeQueryRpcCtx(pJob, pTask, &rpcCtx));
553,752,618✔
1267

1268
      SSubQueryMsg qMsg;
416,760,921✔
1269
      qMsg.header.vgId = addr->nodeId;
553,785,594✔
1270
      qMsg.header.contLen = 0;
553,801,860✔
1271
      qMsg.sId = pTask->seriesId;
553,801,860✔
1272
      qMsg.queryId = pJob->queryId;
553,778,706✔
1273
      qMsg.clientId = pTask->clientId;
553,817,947✔
1274
      qMsg.taskId = pTask->taskId;
553,794,571✔
1275
      qMsg.refId = pJob->refId;
553,800,981✔
1276
      qMsg.execId = pTask->execId;
553,772,455✔
1277
      qMsg.msgMask = (pTask->plan->showRewrite) ? QUERY_MSG_MASK_SHOW_REWRITE() : 0;
553,774,857✔
1278
      qMsg.msgMask |= (pTask->plan->isView) ? QUERY_MSG_MASK_VIEW() : 0;
553,810,375✔
1279
      qMsg.msgMask |= (pTask->plan->isAudit) ? QUERY_MSG_MASK_AUDIT() : 0;
553,798,068✔
1280
      qMsg.msgMask |= (!SCH_IS_PARENT_JOB(pJob) && SCH_IS_ROOT_TASK(pTask)) ? QUERY_MSG_MASK_SUBQUERY() : 0;
553,769,583✔
1281
      qMsg.subQType = (!SCH_IS_PARENT_JOB(pJob) && SCH_IS_ROOT_TASK(pTask)) ? pJob->pDag->subQType : 0;
553,763,394✔
1282
      qMsg.taskType = (pJob->attr.type == JOB_TYPE_HQUERY)? TASK_TYPE_HQUERY:TASK_TYPE_QUERY;
553,784,365✔
1283
      qMsg.explain = SCH_IS_EXPLAIN_JOB(pJob);
553,789,186✔
1284
      qMsg.needFetch = SCH_TASK_NEED_FETCH(pTask);
553,770,740✔
1285
      qMsg.sqlLen = pJob->sql ? strlen(pJob->sql) : 0;
553,786,538✔
1286
      qMsg.sql = pJob->sql;
553,808,745✔
1287
      qMsg.msgLen = pTask->msgLen;
553,797,260✔
1288
      qMsg.msg = pTask->msg;
553,808,075✔
1289

1290
      if (strcmp(tsLocalFqdn, GET_ACTIVE_EP(&addr->epSet)->fqdn) == 0) {
553,784,533✔
1291
        qMsg.compress = 0;
512,055,256✔
1292
      } else {
1293
        qMsg.compress = 1;
41,710,529✔
1294
      }
1295

1296
      SCH_ERR_JRET(schBuildSubJobEndpoints(pTask, &qMsg.subEndPoints, pJob));
553,765,785✔
1297

1298
      msgSize = tSerializeSSubQueryMsg(NULL, 0, &qMsg);
553,704,860✔
1299
      if (msgSize < 0) {
553,500,158✔
UNCOV
1300
        SCH_TASK_ELOG("tSerializeSSubQueryMsg get size, msgSize:%d", msgSize);
×
UNCOV
1301
        taosArrayDestroyP(qMsg.subEndPoints, (FDelete)nodesDestroyNode);
×
UNCOV
1302
        SCH_ERR_JRET(TSDB_CODE_OUT_OF_MEMORY);
×
1303
      }
1304
      
1305
      msg = taosMemoryCalloc(1, msgSize);
553,500,158✔
1306
      if (NULL == msg) {
553,377,898✔
UNCOV
1307
        SCH_TASK_ELOG("calloc %d failed", msgSize);
×
UNCOV
1308
        taosArrayDestroyP(qMsg.subEndPoints, (FDelete)nodesDestroyNode);
×
UNCOV
1309
        SCH_ERR_JRET(terrno);
×
1310
      }
1311

1312
      if (tSerializeSSubQueryMsg(msg, msgSize, &qMsg) < 0) {
553,377,898✔
UNCOV
1313
        SCH_TASK_ELOG("tSerializeSSubQueryMsg failed, msgSize:%d", msgSize);
×
UNCOV
1314
        taosArrayDestroyP(qMsg.subEndPoints, (FDelete)nodesDestroyNode);
×
UNCOV
1315
        taosMemoryFree(msg);
×
1316
        SCH_ERR_JRET(TSDB_CODE_OUT_OF_MEMORY);
×
1317
      }
1318

1319
      taosArrayDestroyP(qMsg.subEndPoints, (FDelete)nodesDestroyNode);
553,201,191✔
1320

1321
      persistHandle = true;
553,435,962✔
1322
      int64_t refId = 0;
553,435,962✔
1323
      code = rpcAllocHandle(&refId);
553,368,545✔
1324
      if (code != 0) {
553,682,519✔
UNCOV
1325
        SCH_TASK_ELOG("rpcAllocHandle failed, code:%x", code);
×
UNCOV
1326
        SCH_ERR_JRET(code);
×
1327
      }
1328

1329
      SCH_SET_TASK_HANDLE(pTask, (void *)refId);
553,527,281✔
1330
      break;
553,573,095✔
1331
    }
1332
    case TDMT_SCH_FETCH:
280,591,138✔
1333
    case TDMT_SCH_MERGE_FETCH: {
1334
      SResFetchReq req = {0};
280,591,138✔
1335
      req.header.vgId = addr->nodeId;
280,591,138✔
1336
      req.sId = pTask->seriesId;
280,591,141✔
1337
      req.queryId = pJob->queryId;
280,590,278✔
1338
      req.clientId = pTask->clientId;
280,591,104✔
1339
      req.taskId = pTask->taskId;
280,590,244✔
1340
      req.execId = pTask->execId;
280,590,904✔
1341
      req.reset = true;
280,591,257✔
1342

1343
      msgSize = tSerializeSResFetchReq(NULL, 0, &req, false, false);
280,591,257✔
1344
      if (msgSize < 0) {
280,590,505✔
UNCOV
1345
        SCH_TASK_ELOG("tSerializeSResFetchReq get size, msgSize:%d", msgSize);
×
UNCOV
1346
        SCH_ERR_RET(TSDB_CODE_OUT_OF_MEMORY);
×
1347
      }
1348
      
1349
      msg = taosMemoryCalloc(1, msgSize);
280,590,505✔
1350
      if (NULL == msg) {
280,590,589✔
UNCOV
1351
        SCH_TASK_ELOG("calloc %d failed", msgSize);
×
UNCOV
1352
        SCH_ERR_RET(terrno);
×
1353
      }
1354

1355
      if (tSerializeSResFetchReq(msg, msgSize, &req, false, false) < 0) {
280,590,589✔
UNCOV
1356
        SCH_TASK_ELOG("tSerializeSResFetchReq %d failed", msgSize);
×
1357
        SCH_ERR_RET(TSDB_CODE_OUT_OF_MEMORY);
8✔
1358
      }
1359
      break;
280,589,602✔
1360
    }
1361
    case TDMT_SCH_DROP_TASK: {
552,061,724✔
1362
      STaskDropReq qMsg;
415,072,632✔
1363
      qMsg.header.vgId = addr->nodeId;
552,061,763✔
1364
      qMsg.header.contLen = 0;
552,061,800✔
1365
      qMsg.sId = pTask->seriesId;
552,061,800✔
1366
      qMsg.queryId = pJob->queryId;
552,060,739✔
1367
      qMsg.clientId = pTask->clientId;
552,060,892✔
1368
      qMsg.taskId = pTask->taskId;
552,061,596✔
1369
      qMsg.refId = pJob->refId;
552,061,477✔
1370
      qMsg.execId = *(int32_t*)param;
552,061,319✔
1371

1372
      msgSize = tSerializeSTaskDropReq(NULL, 0, &qMsg);
552,061,092✔
1373
      if (msgSize < 0) {
552,060,072✔
UNCOV
1374
        SCH_TASK_ELOG("tSerializeSTaskDropReq get size, msgSize:%d", msgSize);
×
UNCOV
1375
        SCH_ERR_RET(TSDB_CODE_OUT_OF_MEMORY);
×
1376
      }
1377
      
1378
      msg = taosMemoryCalloc(1, msgSize);
552,060,072✔
1379
      if (NULL == msg) {
552,062,965✔
UNCOV
1380
        SCH_TASK_ELOG("calloc %d failed", msgSize);
×
UNCOV
1381
        SCH_ERR_RET(terrno);
×
1382
      }
1383

1384
      if (tSerializeSTaskDropReq(msg, msgSize, &qMsg) < 0) {
552,062,965✔
UNCOV
1385
        SCH_TASK_ELOG("tSerializeSTaskDropReq failed, msgSize:%d", msgSize);
×
UNCOV
1386
        taosMemoryFree(msg);
×
1387
        SCH_ERR_RET(TSDB_CODE_OUT_OF_MEMORY);
72✔
1388
      }
1389
      break;
552,060,227✔
1390
    }
1391
/*
1392
    case TDMT_SCH_QUERY_HEARTBEAT: {
1393
      SCH_ERR_RET(schMakeHbRpcCtx(pJob, pTask, &rpcCtx));
1394

1395
      SSchedulerHbReq req = {0};
1396
      req.clientId = schMgmt.clientId;
1397
      req.header.vgId = addr->nodeId;
1398
      req.epId.nodeId = addr->nodeId;
1399
      TAOS_MEMCPY(&req.epId.ep, SCH_GET_CUR_EP(addr), sizeof(SEp));
1400

1401
      msgSize = tSerializeSSchedulerHbReq(NULL, 0, &req);
1402
      if (msgSize < 0) {
1403
        SCH_JOB_ELOG("tSerializeSSchedulerHbReq hbReq failed, size:%d", msgSize);
1404
        SCH_ERR_RET(TSDB_CODE_OUT_OF_MEMORY);
1405
      }
1406
      msg = taosMemoryCalloc(1, msgSize);
1407
      if (NULL == msg) {
1408
        SCH_JOB_ELOG("calloc %d failed", msgSize);
1409
        SCH_ERR_RET(terrno);
1410
      }
1411
      if (tSerializeSSchedulerHbReq(msg, msgSize, &req) < 0) {
1412
        SCH_JOB_ELOG("tSerializeSSchedulerHbReq hbReq failed, size:%d", msgSize);
1413
        SCH_ERR_JRET(TSDB_CODE_OUT_OF_MEMORY);
1414
      }
1415

1416
      persistHandle = true;
1417
      break;
1418
    }
1419
*/    
1420
    case TDMT_SCH_TASK_NOTIFY: {
31,918✔
1421
      ETaskNotifyType* pType = param;
31,918✔
1422
      STaskNotifyReq qMsg = {0};
31,918✔
1423
      qMsg.header.vgId = addr->nodeId;
31,918✔
1424
      qMsg.header.contLen = 0;
31,918✔
1425
      qMsg.sId = pTask->seriesId;
31,918✔
1426
      qMsg.queryId = pJob->queryId;
31,918✔
1427
      qMsg.clientId = pTask->clientId;
31,918✔
1428
      qMsg.taskId = pTask->taskId;
31,918✔
1429
      qMsg.refId = pJob->refId;
31,918✔
1430
      qMsg.execId = pTask->execId;
31,918✔
1431
      qMsg.type = *pType;
31,918✔
1432

1433
      msgSize = tSerializeSTaskNotifyReq(NULL, 0, &qMsg);
31,918✔
1434
      if (msgSize < 0) {
31,918✔
UNCOV
1435
        SCH_TASK_ELOG("tSerializeSTaskNotifyReq get size, msgSize:%d", msgSize);
×
UNCOV
1436
        SCH_ERR_RET(TSDB_CODE_OUT_OF_MEMORY);
×
1437
      }
1438
      
1439
      msg = taosMemoryCalloc(1, msgSize);
31,918✔
1440
      if (NULL == msg) {
31,918✔
UNCOV
1441
        SCH_TASK_ELOG("calloc %d failed", msgSize);
×
UNCOV
1442
        SCH_ERR_RET(terrno);
×
1443
      }
1444

1445
      if (tSerializeSTaskNotifyReq(msg, msgSize, &qMsg) < 0) {
31,918✔
UNCOV
1446
        SCH_TASK_ELOG("tSerializeSTaskNotifyReq failed, msgSize:%d", msgSize);
×
UNCOV
1447
        taosMemoryFree(msg);
×
UNCOV
1448
        SCH_ERR_RET(TSDB_CODE_OUT_OF_MEMORY);
×
1449
      }
1450
      break;      
31,918✔
1451
    }
UNCOV
1452
    default:
×
UNCOV
1453
      SCH_TASK_ELOG("unknown msg type to send, msgType:%d", msgType);
×
UNCOV
1454
      SCH_ERR_RET(TSDB_CODE_SCH_INTERNAL_ERROR);
×
1455
  }
1456

1457
  if ((tsBypassFlag & TSDB_BYPASS_RB_RPC_SEND_SUBMIT) && (TDMT_VND_SUBMIT == msgType)) {
2,107,644,048✔
1458
    taosMemoryFree(msg);
81✔
1459
    SCH_ERR_RET(schProcessOnTaskSuccess(pJob, pTask));
81✔
1460
  } else {
1461
    if (msgType == TDMT_SCH_QUERY || msgType == TDMT_SCH_MERGE_QUERY) {
2,107,643,967✔
1462
      SCH_ERR_JRET(schAppendTaskExecNode(pJob, pTask, addr, pTask->execId));
553,540,815✔
1463
    }
1464

1465
    SSchTrans trans = {.pTrans = pJob->conn.pTrans, .pHandle = SCH_GET_TASK_HANDLE(pTask)};
2,107,867,749✔
1466
    if (msgType == TDMT_SCH_DROP_TASK && pJob->errCode == TSDB_CODE_RPC_TIMEOUT) {
2,107,937,466✔
UNCOV
1467
      trans.pHandle = NULL;
×
UNCOV
1468
      SCH_TASK_WLOG("clear refId before send drop-task msg, code:%s", tstrerror(pJob->errCode));
×
1469
    }
1470

1471
    code = schAsyncSendMsg(pJob, pTask, &trans, addr, msgType, msg, (uint32_t)msgSize, persistHandle, (rpcCtx.args ? &rpcCtx : NULL));
2,107,796,808✔
1472
    msg = NULL;
2,107,867,673✔
1473
    SCH_ERR_JRET(code);
2,107,867,673✔
1474
  }
1475

1476
  return TSDB_CODE_SUCCESS;
2,107,505,282✔
1477

1478
_return:
379,304✔
1479

1480
  pTask->lastMsgType = -1;
365,629✔
1481
  schFreeRpcCtx(&rpcCtx);
365,629✔
1482

1483
  taosMemoryFreeClear(msg);
365,629✔
1484
  SCH_RET(code);
365,629✔
1485
}
1486
// clang-format on
STATUS · Troubleshooting · Open an Issue · Sales · Support · CAREERS · ENTERPRISE · START FREE TRIAL · SCHEDULE DEMO
ANNOUNCEMENTS · TWITTER · TOS & SLA · Supported CI Services · What's a CI service? · Automated Testing

© 2026 Coveralls, Inc