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

taosdata / TDengine / #4897

25 Dec 2025 10:17AM UTC coverage: 65.717% (-0.2%) from 65.929%
#4897

push

travis-ci

web-flow
fix: [6622889291] Fix invalid rowSize. (#34043)

186011 of 283047 relevant lines covered (65.72%)

113853896.64 hits per line

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

75.7
/source/libs/executor/src/executil.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 "function.h"
17
#include "functionMgt.h"
18
#include "index.h"
19
#include "os.h"
20
#include "query.h"
21
#include "querynodes.h"
22
#include "taoserror.h"
23
#include "tarray.h"
24
#include "tcompare.h"
25
#include "tdatablock.h"
26
#include "thash.h"
27
#include "tmsg.h"
28
#include "ttime.h"
29

30
#include "executil.h"
31
#include "executorInt.h"
32
#include "querytask.h"
33
#include "storageapi.h"
34
#include "tutil.h"
35
#include "tjson.h"
36
#include "trpc.h"
37
#include "filter.h"
38

39
typedef struct tagFilterAssist {
40
  SHashObj* colHash;
41
  int32_t   index;
42
  SArray*   cInfoList;
43
  int32_t   code;
44
} tagFilterAssist;
45

46
typedef struct STransTagExprCtx {
47
  int32_t      code;
48
  SMetaReader* pReader;
49
} STransTagExprCtx;
50

51
typedef enum {
52
  FILTER_NO_LOGIC = 1,
53
  FILTER_AND,
54
  FILTER_OTHER,
55
} FilterCondType;
56

57
static FilterCondType checkTagCond(SNode* cond);
58
static int32_t optimizeTbnameInCond(void* metaHandle, int64_t suid, SArray* list, SNode* pTagCond, SStorageAPI* pAPI);
59
static int32_t optimizeTbnameInCondImpl(void* metaHandle, SArray* list, SNode* pTagCond, SStorageAPI* pStoreAPI,
60
                                        uint64_t suid);
61

62
static int32_t getTableList(void* pVnode, SScanPhysiNode* pScanNode, SNode* pTagCond, SNode* pTagIndexCond,
63
                            STableListInfo* pListInfo, uint8_t* digest, const char* idstr, SStorageAPI* pStorageAPI, void* pStreamInfo);
64

65
static int64_t getLimit(const SNode* pLimit) {
948,877,220✔
66
  return (NULL == pLimit || NULL == ((SLimitNode*)pLimit)->limit) ? -1 : ((SLimitNode*)pLimit)->limit->datum.i;
948,877,220✔
67
}
68
static int64_t getOffset(const SNode* pLimit) {
948,794,657✔
69
  return (NULL == pLimit || NULL == ((SLimitNode*)pLimit)->offset) ? -1 : ((SLimitNode*)pLimit)->offset->datum.i;
948,794,657✔
70
}
71
static void releaseColInfoData(void* pCol);
72

73
void initResultRowInfo(SResultRowInfo* pResultRowInfo) {
404,509,995✔
74
  pResultRowInfo->size = 0;
404,509,995✔
75
  pResultRowInfo->cur.pageId = -1;
404,549,328✔
76
}
404,613,959✔
77

78
void closeResultRow(SResultRow* pResultRow) { pResultRow->closed = true; }
6,127,744✔
79

80
void resetResultRow(SResultRow* pResultRow, size_t entrySize) {
399,817,335✔
81
  pResultRow->numOfRows = 0;
399,817,335✔
82
  pResultRow->closed = false;
399,829,931✔
83
  pResultRow->endInterp = false;
399,818,839✔
84
  pResultRow->startInterp = false;
399,830,683✔
85

86
  if (entrySize > 0) {
399,831,999✔
87
    memset(pResultRow->pEntryInfo, 0, entrySize);
399,835,007✔
88
  }
89
}
399,817,899✔
90

91
// TODO refactor: use macro
92
SResultRowEntryInfo* getResultEntryInfo(const SResultRow* pRow, int32_t index, const int32_t* offset) {
2,147,483,647✔
93
  return (SResultRowEntryInfo*)((char*)pRow->pEntryInfo + offset[index]);
2,147,483,647✔
94
}
95

96
size_t getResultRowSize(SqlFunctionCtx* pCtx, int32_t numOfOutput) {
227,863,848✔
97
  int32_t rowSize = (numOfOutput * sizeof(SResultRowEntryInfo)) + sizeof(SResultRow);
227,863,848✔
98

99
  for (int32_t i = 0; i < numOfOutput; ++i) {
901,966,883✔
100
    rowSize += pCtx[i].resDataInfo.interBufSize;
674,193,846✔
101
  }
102

103
  return rowSize;
227,773,037✔
104
}
105

106
// Convert buf read from rocksdb to result row
107
int32_t getResultRowFromBuf(SExprSupp* pSup, const char* inBuf, size_t inBufSize, char** outBuf, size_t* outBufSize) {
×
108
  if (inBuf == NULL || pSup == NULL) {
×
109
    qError("invalid input parameters, inBuf:%p, pSup:%p", inBuf, pSup);
×
110
    return TSDB_CODE_INVALID_PARA;
×
111
  }
112
  SqlFunctionCtx* pCtx = pSup->pCtx;
×
113
  int32_t*        offset = pSup->rowEntryInfoOffset;
×
114
  SResultRow*     pResultRow = NULL;
×
115
  size_t          processedSize = 0;
×
116
  int32_t         code = TSDB_CODE_SUCCESS;
×
117

118
  // calculate the size of output buffer
119
  *outBufSize = getResultRowSize(pCtx, pSup->numOfExprs);
×
120
  *outBuf = taosMemoryMalloc(*outBufSize);
×
121
  if (*outBuf == NULL) {
×
122
    qError("failed to allocate memory for output buffer, size:%zu", *outBufSize);
×
123
    return terrno;
×
124
  }
125
  pResultRow = (SResultRow*)*outBuf;
×
126
  (void)memcpy(pResultRow, inBuf, sizeof(SResultRow));
×
127
  inBuf += sizeof(SResultRow);
×
128
  processedSize += sizeof(SResultRow);
×
129

130
  for (int32_t i = 0; i < pSup->numOfExprs; ++i) {
×
131
    int32_t len = *(int32_t*)inBuf;
×
132
    inBuf += sizeof(int32_t);
×
133
    processedSize += sizeof(int32_t);
×
134
    if (pResultRow->version != FUNCTION_RESULT_INFO_VERSION && pCtx->fpSet.decode) {
×
135
      code = pCtx->fpSet.decode(&pCtx[i], inBuf, getResultEntryInfo(pResultRow, i, offset), pResultRow->version);
×
136
      if (code != TSDB_CODE_SUCCESS) {
×
137
        qError("failed to decode result row, code:%d", code);
×
138
        return code;
×
139
      }
140
    } else {
141
      (void)memcpy(getResultEntryInfo(pResultRow, i, offset), inBuf, len);
×
142
    }
143
    inBuf += len;
×
144
    processedSize += len;
×
145
  }
146

147
  if (processedSize < inBufSize) {
×
148
    // stream stores extra data after result row
149
    size_t leftLen = inBufSize - processedSize;
×
150
    TAOS_MEMORY_REALLOC(*outBuf, *outBufSize + leftLen);
×
151
    if (*outBuf == NULL) {
×
152
      qError("failed to reallocate memory for output buffer, size:%zu", *outBufSize + leftLen);
×
153
      return terrno;
×
154
    }
155
    (void)memcpy(*outBuf + *outBufSize, inBuf, leftLen);
×
156
    inBuf += leftLen;
×
157
    processedSize += leftLen;
×
158
    *outBufSize += leftLen;
×
159
  }
160

161
  qTrace("[StreamInternal] get result inBufSize:%zu, outBufSize:%zu", inBufSize, *outBufSize);
×
162
  return TSDB_CODE_SUCCESS;
×
163
}
164

165
// Convert result row to buf for rocksdb
166
int32_t putResultRowToBuf(SExprSupp* pSup, const char* inBuf, size_t inBufSize, char** outBuf, size_t* outBufSize) {
×
167
  if (pSup == NULL || inBuf == NULL || outBuf == NULL || outBufSize == NULL) {
×
168
    qError("invalid input parameters, inBuf:%p, pSup:%p, outBufSize:%p, outBuf:%p", inBuf, pSup, outBufSize, outBuf);
×
169
    return TSDB_CODE_INVALID_PARA;
×
170
  }
171

172
  SqlFunctionCtx* pCtx = pSup->pCtx;
×
173
  int32_t*        offset = pSup->rowEntryInfoOffset;
×
174
  SResultRow*     pResultRow = (SResultRow*)inBuf;
×
175
  size_t          rowSize = getResultRowSize(pCtx, pSup->numOfExprs);
×
176

177
  if (rowSize > inBufSize) {
×
178
    qError("invalid input buffer size, rowSize:%zu, inBufSize:%zu", rowSize, inBufSize);
×
179
    return TSDB_CODE_INVALID_PARA;
×
180
  }
181

182
  // calculate the size of output buffer
183
  *outBufSize = rowSize + sizeof(int32_t) * pSup->numOfExprs;
×
184
  if (rowSize < inBufSize) {
×
185
    *outBufSize += inBufSize - rowSize;
×
186
  }
187

188
  *outBuf = taosMemoryMalloc(*outBufSize);
×
189
  if (*outBuf == NULL) {
×
190
    qError("failed to allocate memory for output buffer, size:%zu", *outBufSize);
×
191
    return terrno;
×
192
  }
193

194
  char* pBuf = *outBuf;
×
195
  pResultRow->version = FUNCTION_RESULT_INFO_VERSION;
×
196
  (void)memcpy(pBuf, pResultRow, sizeof(SResultRow));
×
197
  pBuf += sizeof(SResultRow);
×
198
  for (int32_t i = 0; i < pSup->numOfExprs; ++i) {
×
199
    size_t len = sizeof(SResultRowEntryInfo) + pCtx[i].resDataInfo.interBufSize;
×
200
    *(int32_t*)pBuf = (int32_t)len;
×
201
    pBuf += sizeof(int32_t);
×
202
    (void)memcpy(pBuf, getResultEntryInfo(pResultRow, i, offset), len);
×
203
    pBuf += len;
×
204
  }
205

206
  if (rowSize < inBufSize) {
×
207
    // stream stores extra data after result row
208
    size_t leftLen = inBufSize - rowSize;
×
209
    (void)memcpy(pBuf, inBuf + rowSize, leftLen);
×
210
    pBuf += leftLen;
×
211
  }
212

213
  qTrace("[StreamInternal] put result inBufSize:%zu, outBufSize:%zu", inBufSize, *outBufSize);
×
214
  return TSDB_CODE_SUCCESS;
×
215
}
216

217
static void freeEx(void* p) { taosMemoryFree(*(void**)p); }
×
218

219
void cleanupGroupResInfo(SGroupResInfo* pGroupResInfo) {
106,004,771✔
220
  taosMemoryFreeClear(pGroupResInfo->pBuf);
106,004,771✔
221
  if (pGroupResInfo->freeItem) {
106,005,361✔
222
    //    taosArrayDestroy(pGroupResInfo->pRows);
223
    taosArrayDestroyEx(pGroupResInfo->pRows, freeEx);
×
224
    pGroupResInfo->freeItem = false;
×
225
    pGroupResInfo->pRows = NULL;
×
226
  } else {
227
    taosArrayDestroy(pGroupResInfo->pRows);
106,003,019✔
228
    pGroupResInfo->pRows = NULL;
106,000,987✔
229
  }
230
  pGroupResInfo->index = 0;
106,000,570✔
231
  pGroupResInfo->delIndex = 0;
106,001,170✔
232
}
106,001,739✔
233

234
int32_t resultrowComparAsc(const void* p1, const void* p2) {
2,147,483,647✔
235
  SResKeyPos* pp1 = *(SResKeyPos**)p1;
2,147,483,647✔
236
  SResKeyPos* pp2 = *(SResKeyPos**)p2;
2,147,483,647✔
237

238
  if (pp1->groupId == pp2->groupId) {
2,147,483,647✔
239
    int64_t pts1 = *(int64_t*)pp1->key;
2,147,483,647✔
240
    int64_t pts2 = *(int64_t*)pp2->key;
2,147,483,647✔
241

242
    if (pts1 == pts2) {
2,147,483,647✔
243
      return 0;
×
244
    } else {
245
      return pts1 < pts2 ? -1 : 1;
2,147,483,647✔
246
    }
247
  } else {
248
    return pp1->groupId < pp2->groupId ? -1 : 1;
2,147,483,647✔
249
  }
250
}
251

252
static int32_t resultrowComparDesc(const void* p1, const void* p2) { return resultrowComparAsc(p2, p1); }
1,943,345,427✔
253

254
int32_t initGroupedResultInfo(SGroupResInfo* pGroupResInfo, SSHashObj* pHashmap, int32_t order) {
70,235,392✔
255
  int32_t code = TSDB_CODE_SUCCESS;
70,235,392✔
256
  int32_t lino = 0;
70,235,392✔
257
  if (pGroupResInfo->pRows != NULL) {
70,235,392✔
258
    taosArrayDestroy(pGroupResInfo->pRows);
5,525,428✔
259
  }
260
  if (pGroupResInfo->pBuf) {
70,233,397✔
261
    taosMemoryFree(pGroupResInfo->pBuf);
5,525,428✔
262
    pGroupResInfo->pBuf = NULL;
5,524,878✔
263
  }
264

265
  // extract the result rows information from the hash map
266
  int32_t size = tSimpleHashGetSize(pHashmap);
70,232,805✔
267

268
  void* pData = NULL;
70,234,465✔
269
  pGroupResInfo->pRows = taosArrayInit(size, POINTER_BYTES);
70,234,465✔
270
  QUERY_CHECK_NULL(pGroupResInfo->pRows, code, lino, _end, terrno);
70,229,215✔
271

272
  size_t  keyLen = 0;
70,230,009✔
273
  int32_t iter = 0;
70,230,637✔
274
  int64_t bufLen = 0, offset = 0;
70,230,215✔
275

276
  // todo move away and record this during create window
277
  while ((pData = tSimpleHashIterate(pHashmap, pData, &iter)) != NULL) {
2,147,483,647✔
278
    /*void* key = */ (void)tSimpleHashGetKey(pData, &keyLen);
279
    bufLen += keyLen + sizeof(SResultRowPosition);
2,147,483,647✔
280
  }
281

282
  pGroupResInfo->pBuf = taosMemoryMalloc(bufLen);
70,223,257✔
283
  QUERY_CHECK_NULL(pGroupResInfo->pBuf, code, lino, _end, terrno);
70,230,448✔
284

285
  iter = 0;
70,224,778✔
286
  while ((pData = tSimpleHashIterate(pHashmap, pData, &iter)) != NULL) {
2,147,483,647✔
287
    void* key = tSimpleHashGetKey(pData, &keyLen);
2,147,483,647✔
288

289
    SResKeyPos* p = (SResKeyPos*)(pGroupResInfo->pBuf + offset);
2,147,483,647✔
290

291
    p->groupId = *(uint64_t*)key;
2,147,483,647✔
292
    p->pos = *(SResultRowPosition*)pData;
2,147,483,647✔
293
    memcpy(p->key, (char*)key + sizeof(uint64_t), keyLen - sizeof(uint64_t));
2,147,483,647✔
294
    void* tmp = taosArrayPush(pGroupResInfo->pRows, &p);
2,147,483,647✔
295
    QUERY_CHECK_NULL(pGroupResInfo->pBuf, code, lino, _end, terrno);
2,147,483,647✔
296

297
    offset += keyLen + sizeof(struct SResultRowPosition);
2,147,483,647✔
298
  }
299

300
  if (order == TSDB_ORDER_ASC || order == TSDB_ORDER_DESC) {
70,215,539✔
301
    __compar_fn_t fn = (order == TSDB_ORDER_ASC) ? resultrowComparAsc : resultrowComparDesc;
13,903,796✔
302
    size = POINTER_BYTES;
13,903,796✔
303
    taosSort(pGroupResInfo->pRows->pData, taosArrayGetSize(pGroupResInfo->pRows), size, fn);
13,903,796✔
304
  }
305

306
  pGroupResInfo->index = 0;
70,216,998✔
307

308
_end:
70,235,663✔
309
  if (code != TSDB_CODE_SUCCESS) {
70,238,778✔
310
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
311
  }
312
  return code;
70,238,778✔
313
}
314

315
void initMultiResInfoFromArrayList(SGroupResInfo* pGroupResInfo, SArray* pArrayList) {
×
316
  if (pGroupResInfo->pRows != NULL) {
×
317
    taosArrayDestroy(pGroupResInfo->pRows);
×
318
  }
319

320
  pGroupResInfo->freeItem = true;
×
321
  pGroupResInfo->pRows = pArrayList;
×
322
  pGroupResInfo->index = 0;
×
323
  pGroupResInfo->delIndex = 0;
×
324
}
×
325

326
bool hasRemainResults(SGroupResInfo* pGroupResInfo) {
240,307,835✔
327
  if (pGroupResInfo->pRows == NULL) {
240,307,835✔
328
    return false;
×
329
  }
330

331
  return pGroupResInfo->index < taosArrayGetSize(pGroupResInfo->pRows);
240,312,565✔
332
}
333

334
int32_t getNumOfTotalRes(SGroupResInfo* pGroupResInfo) {
129,060,097✔
335
  if (pGroupResInfo->pRows == 0) {
129,060,097✔
336
    return 0;
×
337
  }
338

339
  return (int32_t)taosArrayGetSize(pGroupResInfo->pRows);
129,064,082✔
340
}
341

342
SArray* createSortInfo(SNodeList* pNodeList) {
42,584,989✔
343
  size_t numOfCols = 0;
42,584,989✔
344

345
  if (pNodeList != NULL) {
42,584,989✔
346
    numOfCols = LIST_LENGTH(pNodeList);
42,531,752✔
347
  } else {
348
    numOfCols = 0;
53,310✔
349
  }
350

351
  SArray* pList = taosArrayInit(numOfCols, sizeof(SBlockOrderInfo));
42,586,868✔
352
  if (pList == NULL) {
42,581,094✔
353
    return pList;
×
354
  }
355

356
  for (int32_t i = 0; i < numOfCols; ++i) {
94,711,088✔
357
    SOrderByExprNode* pSortKey = (SOrderByExprNode*)nodesListGetNode(pNodeList, i);
52,127,380✔
358
    if (!pSortKey) {
52,128,807✔
359
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
360
      taosArrayDestroy(pList);
×
361
      pList = NULL;
×
362
      terrno = TSDB_CODE_QRY_EXECUTOR_INTERNAL_ERROR;
×
363
      break;
×
364
    }
365
    SBlockOrderInfo bi = {0};
52,128,807✔
366
    bi.order = (pSortKey->order == ORDER_ASC) ? TSDB_ORDER_ASC : TSDB_ORDER_DESC;
52,127,095✔
367
    bi.nullFirst = (pSortKey->nullOrder == NULL_ORDER_FIRST);
52,128,249✔
368

369
    if (nodeType(pSortKey->pExpr) != QUERY_NODE_COLUMN) {
52,129,655✔
370
      qError("invalid order by expr type:%d", nodeType(pSortKey->pExpr));
×
371
      taosArrayDestroy(pList);
×
372
      pList = NULL;
×
373
      terrno = TSDB_CODE_QRY_EXECUTOR_INTERNAL_ERROR;
×
374
      break;
×
375
    }
376
    
377
    SColumnNode* pColNode = (SColumnNode*)pSortKey->pExpr;
52,121,681✔
378
    bi.slotId = pColNode->slotId;
52,128,893✔
379
    void* tmp = taosArrayPush(pList, &bi);
52,131,507✔
380
    if (!tmp) {
52,131,507✔
381
      taosArrayDestroy(pList);
×
382
      pList = NULL;
×
383
      break;
×
384
    }
385
  }
386

387
  return pList;
42,584,524✔
388
}
389

390
SSDataBlock* createDataBlockFromDescNode(void* p) {
586,135,647✔
391
  SDataBlockDescNode* pNode = (SDataBlockDescNode*)p;
586,135,647✔
392
  int32_t      numOfCols = LIST_LENGTH(pNode->pSlots);
586,135,647✔
393
  SSDataBlock* pBlock = NULL;
586,302,073✔
394
  int32_t      code = createDataBlock(&pBlock);
586,296,059✔
395
  if (code) {
585,966,816✔
396
    terrno = code;
×
397
    return NULL;
×
398
  }
399

400
  pBlock->info.id.blockId = pNode->dataBlockId;
585,966,816✔
401
  pBlock->info.type = STREAM_INVALID;
585,988,622✔
402
  pBlock->info.calWin = (STimeWindow){.skey = INT64_MIN, .ekey = INT64_MAX};
586,034,214✔
403
  pBlock->info.watermark = INT64_MIN;
586,214,883✔
404

405
  for (int32_t i = 0; i < numOfCols; ++i) {
2,147,483,647✔
406
    SSlotDescNode* pDescNode = (SSlotDescNode*)nodesListGetNode(pNode->pSlots, i);
2,128,616,158✔
407
    if (!pDescNode) {
2,128,515,048✔
408
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(code));
×
409
      blockDataDestroy(pBlock);
×
410
      pBlock = NULL;
×
411
      terrno = TSDB_CODE_INVALID_PARA;
×
412
      break;
×
413
    }
414
    SColumnInfoData idata =
2,128,459,306✔
415
        createColumnInfoData(pDescNode->dataType.type, pDescNode->dataType.bytes, pDescNode->slotId);
2,128,677,620✔
416
    idata.info.scale = pDescNode->dataType.scale;
2,128,813,886✔
417
    idata.info.precision = pDescNode->dataType.precision;
2,128,710,246✔
418
    idata.info.noData = pDescNode->reserve;
2,128,857,316✔
419

420
    code = blockDataAppendColInfo(pBlock, &idata);
2,128,792,984✔
421
    if (code != TSDB_CODE_SUCCESS) {
2,128,847,055✔
422
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(code));
115,620✔
423
      blockDataDestroy(pBlock);
115,620✔
424
      pBlock = NULL;
×
425
      terrno = code;
×
426
      break;
×
427
    }
428
  }
429

430
  return pBlock;
586,439,655✔
431
}
432

433
int32_t prepareDataBlockBuf(SSDataBlock* pDataBlock, SColMatchInfo* pMatchInfo) {
211,470,795✔
434
  SDataBlockInfo* pBlockInfo = &pDataBlock->info;
211,470,795✔
435

436
  for (int32_t i = 0; i < taosArrayGetSize(pMatchInfo->pList); ++i) {
990,310,982✔
437
    SColMatchItem* pItem = taosArrayGet(pMatchInfo->pList, i);
786,156,719✔
438
    if (!pItem) {
785,995,824✔
439
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
440
      return terrno;
×
441
    }
442

443
    if (pItem->isPk) {
785,995,824✔
444
      SColumnInfoData* pInfoData = taosArrayGet(pDataBlock->pDataBlock, pItem->dstSlotId);
7,520,686✔
445
      if (!pInfoData) {
6,980,122✔
446
        qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
447
        return terrno;
×
448
      }
449
      pBlockInfo->pks[0].type = pInfoData->info.type;
6,980,122✔
450
      pBlockInfo->pks[1].type = pInfoData->info.type;
6,993,136✔
451

452
      // allocate enough buffer size, which is pInfoData->info.bytes
453
      if (IS_VAR_DATA_TYPE(pItem->dataType.type)) {
6,992,668✔
454
        pBlockInfo->pks[0].pData = taosMemoryCalloc(1, pInfoData->info.bytes);
2,340,834✔
455
        if (pBlockInfo->pks[0].pData == NULL) {
2,331,249✔
456
          return terrno;
×
457
        }
458

459
        pBlockInfo->pks[1].pData = taosMemoryCalloc(1, pInfoData->info.bytes);
2,336,301✔
460
        if (pBlockInfo->pks[1].pData == NULL) {
2,333,340✔
461
          taosMemoryFreeClear(pBlockInfo->pks[0].pData);
×
462
          return terrno;
×
463
        }
464

465
        pBlockInfo->pks[0].nData = pInfoData->info.bytes;
2,337,522✔
466
        pBlockInfo->pks[1].nData = pInfoData->info.bytes;
2,336,825✔
467
      }
468

469
      break;
6,986,863✔
470
    }
471
  }
472

473
  return TSDB_CODE_SUCCESS;
211,210,841✔
474
}
475

476
EDealRes doTranslateTagExpr(SNode** pNode, void* pContext) {
372,414✔
477
  STransTagExprCtx* pCtx = pContext;
372,414✔
478
  SMetaReader*      mr = pCtx->pReader;
372,414✔
479
  bool              isTagCol = false, isTbname = false;
372,414✔
480
  if (nodeType(*pNode) == QUERY_NODE_COLUMN) {
372,414✔
481
    SColumnNode* pCol = (SColumnNode*)*pNode;
106,404✔
482
    if (pCol->colType == COLUMN_TYPE_TBNAME)
106,404✔
483
      isTbname = true;
×
484
    else
485
      isTagCol = true;
106,404✔
486
  } else if (nodeType(*pNode) == QUERY_NODE_FUNCTION) {
266,010✔
487
    SFunctionNode* pFunc = (SFunctionNode*)*pNode;
×
488
    if (pFunc->funcType == FUNCTION_TYPE_TBNAME) isTbname = true;
×
489
  }
490
  if (isTagCol) {
372,414✔
491
    SColumnNode* pSColumnNode = *(SColumnNode**)pNode;
106,404✔
492

493
    SValueNode* res = NULL;
106,404✔
494
    pCtx->code = nodesMakeNode(QUERY_NODE_VALUE, (SNode**)&res);
106,404✔
495
    if (NULL == res) {
106,404✔
496
      return DEAL_RES_ERROR;
×
497
    }
498

499
    res->translate = true;
106,404✔
500
    res->node.resType = pSColumnNode->node.resType;
106,404✔
501

502
    STagVal tagVal = {0};
106,404✔
503
    tagVal.cid = pSColumnNode->colId;
106,404✔
504
    const char* p = mr->pAPI->extractTagVal(mr->me.ctbEntry.pTags, pSColumnNode->node.resType.type, &tagVal);
106,404✔
505
    if (p == NULL) {
106,404✔
506
      res->node.resType.type = TSDB_DATA_TYPE_NULL;
×
507
    } else if (pSColumnNode->node.resType.type == TSDB_DATA_TYPE_JSON) {
106,404✔
508
      int32_t len = ((const STag*)p)->len;
×
509
      res->datum.p = taosMemoryCalloc(len + 1, 1);
×
510
      if (NULL == res->datum.p) {
×
511
        return DEAL_RES_ERROR;
×
512
      }
513
      memcpy(res->datum.p, p, len);
×
514
    } else if (IS_VAR_DATA_TYPE(pSColumnNode->node.resType.type)) {
106,404✔
515
      if (IS_STR_DATA_BLOB(pSColumnNode->node.resType.type)) {
106,404✔
516
        return TSDB_CODE_BLOB_NOT_SUPPORT_TAG;
×
517
      }
518

519
      res->datum.p = taosMemoryCalloc(tagVal.nData + VARSTR_HEADER_SIZE + 1, 1);
106,404✔
520
      if (NULL == res->datum.p) {
106,404✔
521
        return DEAL_RES_ERROR;
×
522
      }
523

524
      if (IS_STR_DATA_BLOB(pSColumnNode->node.resType.type)) {
106,404✔
525
        memcpy(blobDataVal(res->datum.p), tagVal.pData, tagVal.nData);
×
526
        blobDataSetLen(res->datum.p, tagVal.nData);
×
527
      } else {
528
        memcpy(varDataVal(res->datum.p), tagVal.pData, tagVal.nData);
106,404✔
529
        varDataSetLen(res->datum.p, tagVal.nData);
106,404✔
530
      }
531
    } else {
532
      int32_t code = nodesSetValueNodeValue(res, &(tagVal.i64));
×
533
      if (code != TSDB_CODE_SUCCESS) {
×
534
        return DEAL_RES_ERROR;
×
535
      }
536
    }
537
    nodesDestroyNode(*pNode);
106,404✔
538
    *pNode = (SNode*)res;
106,404✔
539
  } else if (isTbname) {
266,010✔
540
    SValueNode* res = NULL;
×
541
    pCtx->code = nodesMakeNode(QUERY_NODE_VALUE, (SNode**)&res);
×
542
    if (NULL == res) {
×
543
      return DEAL_RES_ERROR;
×
544
    }
545

546
    res->translate = true;
×
547
    res->node.resType = ((SExprNode*)(*pNode))->resType;
×
548

549
    int32_t len = strlen(mr->me.name);
×
550
    res->datum.p = taosMemoryCalloc(len + VARSTR_HEADER_SIZE + 1, 1);
×
551
    if (NULL == res->datum.p) {
×
552
      return DEAL_RES_ERROR;
×
553
    }
554
    memcpy(varDataVal(res->datum.p), mr->me.name, len);
×
555
    varDataSetLen(res->datum.p, len);
×
556
    nodesDestroyNode(*pNode);
×
557
    *pNode = (SNode*)res;
×
558
  }
559

560
  return DEAL_RES_CONTINUE;
372,414✔
561
}
562

563
int32_t isQualifiedTable(int64_t uid, SNode* pTagCond, void* vnode, bool* pQualified, SStorageAPI* pAPI) {
53,202✔
564
  int32_t     code = TSDB_CODE_SUCCESS;
53,202✔
565
  SMetaReader mr = {0};
53,202✔
566

567
  pAPI->metaReaderFn.initReader(&mr, vnode, META_READER_LOCK, &pAPI->metaFn);
53,202✔
568
  code = pAPI->metaReaderFn.getEntryGetUidCache(&mr, uid);
53,202✔
569
  if (TSDB_CODE_SUCCESS != code) {
53,202✔
570
    pAPI->metaReaderFn.clearReader(&mr);
×
571
    *pQualified = false;
×
572

573
    return TSDB_CODE_SUCCESS;
×
574
  }
575

576
  SNode* pTagCondTmp = NULL;
53,202✔
577
  code = nodesCloneNode(pTagCond, &pTagCondTmp);
53,202✔
578
  if (TSDB_CODE_SUCCESS != code) {
53,202✔
579
    *pQualified = false;
×
580
    pAPI->metaReaderFn.clearReader(&mr);
×
581
    return code;
×
582
  }
583
  STransTagExprCtx ctx = {.code = 0, .pReader = &mr};
53,202✔
584
  nodesRewriteExprPostOrder(&pTagCondTmp, doTranslateTagExpr, &ctx);
53,202✔
585
  pAPI->metaReaderFn.clearReader(&mr);
53,202✔
586
  if (TSDB_CODE_SUCCESS != ctx.code) {
53,202✔
587
    *pQualified = false;
×
588
    nodesDestroyNode(pTagCondTmp);
×
589
    terrno = code;
×
590
    return code;
×
591
  }
592

593
  SNode* pNew = NULL;
53,202✔
594
  code = scalarCalculateConstants(pTagCondTmp, &pNew);
53,202✔
595
  if (TSDB_CODE_SUCCESS != code) {
53,202✔
596
    terrno = code;
×
597
    nodesDestroyNode(pTagCondTmp);
×
598
    *pQualified = false;
×
599

600
    return code;
×
601
  }
602

603
  SValueNode* pValue = (SValueNode*)pNew;
53,202✔
604
  *pQualified = pValue->datum.b;
53,202✔
605

606
  nodesDestroyNode(pNew);
53,202✔
607
  return TSDB_CODE_SUCCESS;
53,202✔
608
}
609

610
static EDealRes getColumn(SNode** pNode, void* pContext) {
45,844,350✔
611
  tagFilterAssist* pData = (tagFilterAssist*)pContext;
45,844,350✔
612
  SColumnNode*     pSColumnNode = NULL;
45,844,350✔
613
  if (QUERY_NODE_COLUMN == nodeType((*pNode))) {
45,845,712✔
614
    pSColumnNode = *(SColumnNode**)pNode;
14,864,007✔
615
  } else if (QUERY_NODE_FUNCTION == nodeType((*pNode))) {
30,985,246✔
616
    SFunctionNode* pFuncNode = *(SFunctionNode**)(pNode);
697,834✔
617
    if (pFuncNode->funcType == FUNCTION_TYPE_TBNAME) {
697,834✔
618
      pData->code = nodesMakeNode(QUERY_NODE_COLUMN, (SNode**)&pSColumnNode);
646,561✔
619
      if (NULL == pSColumnNode) {
647,288✔
620
        return DEAL_RES_ERROR;
×
621
      }
622
      pSColumnNode->colId = -1;
647,288✔
623
      pSColumnNode->colType = COLUMN_TYPE_TBNAME;
647,288✔
624
      pSColumnNode->node.resType.type = TSDB_DATA_TYPE_VARCHAR;
646,770✔
625
      pSColumnNode->node.resType.bytes = TSDB_TABLE_FNAME_LEN - 1 + VARSTR_HEADER_SIZE;
647,288✔
626
      nodesDestroyNode(*pNode);
647,079✔
627
      *pNode = (SNode*)pSColumnNode;
647,288✔
628
    } else {
629
      return DEAL_RES_CONTINUE;
50,546✔
630
    }
631
  } else {
632
    return DEAL_RES_CONTINUE;
30,287,551✔
633
  }
634

635
  void* data = taosHashGet(pData->colHash, &pSColumnNode->colId, sizeof(pSColumnNode->colId));
15,512,341✔
636
  if (!data) {
15,506,973✔
637
    int32_t tempRes =
638
        taosHashPut(pData->colHash, &pSColumnNode->colId, sizeof(pSColumnNode->colId), pNode, sizeof((*pNode)));
14,177,925✔
639
    if (tempRes != TSDB_CODE_SUCCESS && tempRes != TSDB_CODE_DUP_KEY) {
14,182,436✔
640
      return DEAL_RES_ERROR;
×
641
    }
642
    pSColumnNode->slotId = pData->index++;
14,182,436✔
643
    SColumnInfo cInfo = {.colId = pSColumnNode->colId,
14,182,096✔
644
                         .type = pSColumnNode->node.resType.type,
14,179,648✔
645
                         .bytes = pSColumnNode->node.resType.bytes,
14,182,216✔
646
                         .pk = pSColumnNode->isPk};
14,178,776✔
647
#if TAG_FILTER_DEBUG
648
    qDebug("tagfilter build column info, slotId:%d, colId:%d, type:%d", pSColumnNode->slotId, cInfo.colId, cInfo.type);
649
#endif
650
    void* tmp = taosArrayPush(pData->cInfoList, &cInfo);
14,182,358✔
651
    if (!tmp) {
14,181,936✔
652
      return DEAL_RES_ERROR;
×
653
    }
654
  } else {
655
    SColumnNode* col = *(SColumnNode**)data;
1,329,048✔
656
    pSColumnNode->slotId = col->slotId;
1,329,048✔
657
  }
658

659
  return DEAL_RES_CONTINUE;
15,510,602✔
660
}
661

662
static int32_t createResultData(SDataType* pType, int32_t numOfRows, SScalarParam* pParam) {
13,204,065✔
663
  SColumnInfoData* pColumnData = taosMemoryCalloc(1, sizeof(SColumnInfoData));
13,204,065✔
664
  if (pColumnData == NULL) {
13,205,304✔
665
    return terrno;
×
666
  }
667

668
  pColumnData->info.type = pType->type;
13,205,304✔
669
  pColumnData->info.bytes = pType->bytes;
13,203,659✔
670
  pColumnData->info.scale = pType->scale;
13,204,338✔
671
  pColumnData->info.precision = pType->precision;
13,199,198✔
672

673
  int32_t code = colInfoDataEnsureCapacity(pColumnData, numOfRows, true);
13,203,157✔
674
  if (code != TSDB_CODE_SUCCESS) {
13,198,477✔
675
    terrno = code;
×
676
    releaseColInfoData(pColumnData);
×
677
    return terrno;
×
678
  }
679

680
  pParam->columnData = pColumnData;
13,198,477✔
681
  pParam->colAlloced = true;
13,204,295✔
682
  return TSDB_CODE_SUCCESS;
13,203,443✔
683
}
684

685
static void releaseColInfoData(void* pCol) {
2,055,113✔
686
  if (pCol) {
2,055,113✔
687
    SColumnInfoData* col = (SColumnInfoData*)pCol;
2,055,113✔
688
    colDataDestroy(col);
2,055,113✔
689
    taosMemoryFree(col);
2,054,970✔
690
  }
691
}
2,055,113✔
692

693
void freeItem(void* p) {
189,029,780✔
694
  STUidTagInfo* pInfo = p;
189,029,780✔
695
  if (pInfo->pTagVal != NULL) {
189,029,780✔
696
    taosMemoryFree(pInfo->pTagVal);
188,604,557✔
697
  }
698
}
189,029,859✔
699

700
typedef struct {
701
  col_id_t  colId;
702
  SNode*    pValueNode;
703
  int32_t   bytes;  // length defined in schema
704
} STagDataEntry;
705

706
static int compareTagDataEntry(const void* a, const void* b) {
45,360✔
707
  STagDataEntry* p1 = (STagDataEntry*)a;
45,360✔
708
  STagDataEntry* p2 = (STagDataEntry*)b;
45,360✔
709
  return compareInt16Val(&p1->colId, &p2->colId);
45,360✔
710
}
711

712
static int32_t buildTagDataEntryKey(SArray* pIdWithValue, char** keyBuf, int32_t keyLen) {
22,680✔
713
  *keyBuf = (char*)taosMemoryCalloc(1, keyLen);
22,680✔
714
  if (NULL == *keyBuf) {
22,680✔
715
    qError(
×
716
      "failed to allocate memory for tag filter optimization key, size:%d",
717
      keyLen);
718
    return terrno;
×
719
  }
720
  char* pStart = *keyBuf;
22,680✔
721
  for (int32_t i = 0; i < taosArrayGetSize(pIdWithValue); ++i) {
67,680✔
722
    STagDataEntry* entry      = (STagDataEntry*)taosArrayGet(pIdWithValue, i);
45,000✔
723
    SValueNode*    pValueNode = (SValueNode*)entry->pValueNode;
45,360✔
724
    // num type may have different bytes length, use the smaller one
725
    int32_t        bytes = TMIN(entry->bytes, pValueNode->node.resType.bytes);
45,360✔
726

727
    (void)memcpy(pStart, &entry->colId, sizeof(col_id_t));
45,360✔
728
    pStart += sizeof(col_id_t);
45,360✔
729

730
    if (!pValueNode->isNull) {
45,360✔
731
      switch (pValueNode->node.resType.type) {
39,960✔
732
        case TSDB_DATA_TYPE_BOOL:
4,680✔
733
          (void)memcpy(
4,680✔
734
            pStart, &pValueNode->datum.b, bytes);
4,680✔
735
          pStart += bytes;
4,680✔
736
          break;
4,680✔
737
        case TSDB_DATA_TYPE_TINYINT:
19,440✔
738
        case TSDB_DATA_TYPE_SMALLINT:
739
        case TSDB_DATA_TYPE_INT:
740
        case TSDB_DATA_TYPE_BIGINT:
741
        case TSDB_DATA_TYPE_TIMESTAMP:
742
          (void)memcpy(
19,440✔
743
            pStart, &pValueNode->datum.i, bytes);
19,440✔
744
          pStart += bytes;
19,440✔
745
          break;
19,440✔
746
        case TSDB_DATA_TYPE_UTINYINT:
×
747
        case TSDB_DATA_TYPE_USMALLINT:
748
        case TSDB_DATA_TYPE_UINT:
749
        case TSDB_DATA_TYPE_UBIGINT:
750
          (void)memcpy(
×
751
            pStart, &pValueNode->datum.u, bytes);
×
752
          pStart += bytes;
×
753
          break;
×
754
        case TSDB_DATA_TYPE_FLOAT:
5,400✔
755
        case TSDB_DATA_TYPE_DOUBLE:
756
          (void)memcpy(
5,400✔
757
            pStart, &pValueNode->datum.d, bytes);
5,400✔
758
          pStart += bytes;
5,400✔
759
          break;
5,400✔
760
        case TSDB_DATA_TYPE_VARCHAR:
10,080✔
761
        case TSDB_DATA_TYPE_VARBINARY:
762
        case TSDB_DATA_TYPE_NCHAR:
763
          (void)memcpy(pStart,
10,080✔
764
            varDataVal(pValueNode->datum.p), varDataLen(pValueNode->datum.p));
10,080✔
765
          pStart += varDataLen(pValueNode->datum.p);
10,080✔
766
          break;
10,080✔
767
        default:
×
768
          qError("unsupported tag data type %d in tag filter optimization",
×
769
            pValueNode->node.resType.type);
770
          return TSDB_CODE_STREAM_INTERNAL_ERROR;
×
771
      }
772
    }
773
  }
774

775
  return TSDB_CODE_SUCCESS;
22,680✔
776
}
777

778
static void extractTagDataEntry(
45,360✔
779
  SOperatorNode* pOpNode, SArray* pIdWithValue) {
780
  SNode* pLeft = pOpNode->pLeft;
45,360✔
781
  SNode* pRight = pOpNode->pRight;
45,000✔
782
  SColumnNode* pColNode = nodeType(pLeft) == QUERY_NODE_COLUMN ?
45,360✔
783
    (SColumnNode*)pLeft : (SColumnNode*)pRight;
45,000✔
784
  SValueNode* pValueNode = nodeType(pLeft) == QUERY_NODE_VALUE ?
45,000✔
785
    (SValueNode*)pLeft : (SValueNode*)pRight;
45,000✔
786

787
  STagDataEntry entry = {0};
45,000✔
788
  entry.colId = pColNode->colId;
45,000✔
789
  entry.pValueNode = (SNode*)pValueNode;
45,360✔
790
  entry.bytes = pColNode->node.resType.bytes;
45,360✔
791
  void* _tmp = taosArrayPush(pIdWithValue, &entry);
45,360✔
792
}
45,360✔
793

794
static int32_t extractTagFilterTagDataEntries(
22,680✔
795
  const SNode* pTagCond, SArray* pIdWithVal) {
796
  if (NULL == pTagCond || NULL == pIdWithVal ||
22,680✔
797
    (nodeType(pTagCond) != QUERY_NODE_OPERATOR &&
22,680✔
798
      nodeType(pTagCond) != QUERY_NODE_LOGIC_CONDITION)) {
22,680✔
799
    qError("invalid parameter to extract tag filter symbol");
×
800
    return TSDB_CODE_STREAM_INTERNAL_ERROR;
×
801
  }
802

803
  if (nodeType(pTagCond) == QUERY_NODE_OPERATOR) {
22,680✔
804
    extractTagDataEntry((SOperatorNode*)pTagCond, pIdWithVal);
×
805
  } else if (nodeType(pTagCond) == QUERY_NODE_LOGIC_CONDITION) {
22,680✔
806
    SNode* pChild = NULL;
22,680✔
807
    FOREACH(pChild, ((SLogicConditionNode*)pTagCond)->pParameterList) {
68,040✔
808
      extractTagDataEntry((SOperatorNode*)pChild, pIdWithVal);
45,360✔
809
    }
810
  }
811

812
  taosArraySort(pIdWithVal, compareTagDataEntry);
22,680✔
813

814
  return TSDB_CODE_SUCCESS;
22,680✔
815
}
816

817
static int32_t genStableTagFilterDigest(const SNode* pTagCond, T_MD5_CTX* pContext) {
22,680✔
818
  if (pTagCond == NULL) {
22,680✔
819
    return TSDB_CODE_SUCCESS;
×
820
  }
821

822
  char*   payload = NULL;
22,680✔
823
  int32_t len = 0;
22,680✔
824
  int32_t code = TSDB_CODE_SUCCESS;
22,680✔
825
  int32_t lino = 0;
22,680✔
826

827
  SArray* pIdWithVal = taosArrayInit(TARRAY_MIN_SIZE, sizeof(STagDataEntry));
22,680✔
828
  code = extractTagFilterTagDataEntries(pTagCond, pIdWithVal);
22,320✔
829
  QUERY_CHECK_CODE(code, lino, _end);
22,680✔
830
  for (int32_t i = 0; i < taosArrayGetSize(pIdWithVal); ++i) {
67,680✔
831
    STagDataEntry* pEntry = taosArrayGet(pIdWithVal, i);
45,000✔
832
    len += sizeof(col_id_t) + pEntry->bytes;
45,000✔
833
  }
834
  code = buildTagDataEntryKey(pIdWithVal, &payload, len);
22,320✔
835
  QUERY_CHECK_CODE(code, lino, _end);
22,680✔
836

837
  tMD5Init(pContext);
22,680✔
838
  tMD5Update(pContext, (uint8_t*)payload, (uint32_t)len);
22,680✔
839
  tMD5Final(pContext);
22,680✔
840

841
_end:
22,680✔
842
  if (TSDB_CODE_SUCCESS != code) {
22,680✔
843
    qError("%s failed at line %d since %s",
×
844
      __func__, __LINE__, tstrerror(code));
845
  }
846
  taosArrayDestroy(pIdWithVal);
22,680✔
847
  taosMemoryFree(payload);
22,680✔
848
  return code;
22,680✔
849
}
850

851
static int32_t genTagFilterDigest(const SNode* pTagCond, T_MD5_CTX* pContext) {
63,396✔
852
  if (pTagCond == NULL) {
63,396✔
853
    return TSDB_CODE_SUCCESS;
60,548✔
854
  }
855

856
  char*   payload = NULL;
2,848✔
857
  int32_t len = 0;
2,848✔
858
  int32_t code = nodesNodeToMsg(pTagCond, &payload, &len);
2,848✔
859
  if (code != TSDB_CODE_SUCCESS) {
2,848✔
860
    qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(code));
×
861
    return code;
×
862
  }
863

864
  tMD5Init(pContext);
2,848✔
865
  tMD5Update(pContext, (uint8_t*)payload, (uint32_t)len);
2,848✔
866
  tMD5Final(pContext);
2,848✔
867

868
  // void* tmp = NULL;
869
  // uint32_t size = 0;
870
  // (void)taosAscii2Hex((const char*)pContext->digest, 16, &tmp, &size);
871
  // qInfo("tag filter digest payload: %s", tmp);
872
  // taosMemoryFree(tmp);
873

874
  taosMemoryFree(payload);
2,848✔
875
  return TSDB_CODE_SUCCESS;
2,848✔
876
}
877

878
static int32_t genTbGroupDigest(const SNode* pGroup, uint8_t* filterDigest, T_MD5_CTX* pContext) {
×
879
  int32_t code = TSDB_CODE_SUCCESS;
×
880
  int32_t lino = 0;
×
881
  char*   payload = NULL;
×
882
  int32_t len = 0;
×
883
  code = nodesNodeToMsg(pGroup, &payload, &len);
×
884
  QUERY_CHECK_CODE(code, lino, _end);
×
885

886
  if (filterDigest[0]) {
×
887
    payload = taosMemoryRealloc(payload, len + tListLen(pContext->digest));
×
888
    QUERY_CHECK_NULL(payload, code, lino, _end, terrno);
×
889
    memcpy(payload + len, filterDigest + 1, tListLen(pContext->digest));
×
890
    len += tListLen(pContext->digest);
×
891
  }
892

893
  tMD5Init(pContext);
×
894
  tMD5Update(pContext, (uint8_t*)payload, (uint32_t)len);
×
895
  tMD5Final(pContext);
×
896

897
_end:
×
898
  taosMemoryFree(payload);
×
899
  if (code != TSDB_CODE_SUCCESS) {
×
900
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
901
  }
902
  return code;
×
903
}
904

905
int32_t qGetColumnsFromNodeList(void* data, bool isList, SArray** pColList) {
13,029,940✔
906
  int32_t code = TSDB_CODE_SUCCESS;
13,029,940✔
907
  tagFilterAssist ctx = {0};
13,029,940✔
908
  ctx.colHash = taosHashInit(4, taosGetDefaultHashFunction(TSDB_DATA_TYPE_SMALLINT), false, HASH_NO_LOCK);
13,030,152✔
909
  if (ctx.colHash == NULL) {
13,029,046✔
910
    code = terrno;
×
911
    goto end;
×
912
  }
913

914
  ctx.index = 0;
13,029,046✔
915
  ctx.cInfoList = taosArrayInit(4, sizeof(SColumnInfo));
13,029,046✔
916
  if (ctx.cInfoList == NULL) {
13,029,848✔
917
    code = terrno;
1,629✔
918
    goto end;
×
919
  }
920

921
  if (isList) {
13,028,219✔
922
    SNode* pNode = NULL;
1,877,683✔
923
    FOREACH(pNode, (SNodeList*)data) {
3,940,212✔
924
      nodesRewriteExprPostOrder(&pNode, getColumn, (void*)&ctx);
2,062,529✔
925
      if (TSDB_CODE_SUCCESS != ctx.code) {
2,061,603✔
926
        code = ctx.code;
×
927
        goto end;
×
928
      }
929
      REPLACE_NODE(pNode);
2,061,603✔
930
    }
931
  } else {
932
    SNode* pNode = (SNode*)data;
11,150,536✔
933
    nodesRewriteExprPostOrder(&pNode, getColumn, (void*)&ctx);
11,151,362✔
934
    if (TSDB_CODE_SUCCESS != ctx.code) {
11,152,991✔
935
      code = ctx.code;
×
936
      goto end;
×
937
    }
938
  }
939
  
940
  if (pColList != NULL) *pColList = ctx.cInfoList;
13,028,705✔
941
  ctx.cInfoList = NULL;
13,029,040✔
942

943
end:
13,029,940✔
944
  taosHashCleanup(ctx.colHash);
13,028,734✔
945
  taosArrayDestroy(ctx.cInfoList);
13,017,201✔
946
  return code;
13,018,637✔
947
}
948

949
static int32_t buildGroupInfo(SColumnInfoData* pValue, int32_t i, SArray* gInfo) {
667,134✔
950
  int32_t code = TSDB_CODE_SUCCESS;
667,134✔
951
  SStreamGroupValue* v = taosArrayReserve(gInfo, 1);
667,134✔
952
  if (v == NULL) {
667,134✔
953
    code = terrno;
×
954
    goto end;
×
955
  }
956
  if (colDataIsNull_s(pValue, i)) {
1,334,268✔
957
    v->isNull = true;
14,400✔
958
  } else {
959
    v->isNull = false;
652,734✔
960
    char* data = colDataGetData(pValue, i);
652,734✔
961
    if (pValue->info.type == TSDB_DATA_TYPE_JSON) {
652,374✔
962
      if (tTagIsJson(data)) {
×
963
        code = TSDB_CODE_QRY_JSON_IN_GROUP_ERROR;
×
964
        goto end;
×
965
      }
966
      if (tTagIsJsonNull(data)) {
×
967
        v->isNull = true;
×
968
        goto end;
×
969
      }
970
      int32_t len = getJsonValueLen(data);
×
971
      v->data.type = pValue->info.type;
×
972
      v->data.nData = len;
×
973
      v->data.pData = taosMemoryCalloc(1, len + 1);
×
974
      if (v->data.pData == NULL) {
×
975
        code = terrno;
×
976
        goto end;
×
977
      }
978
      memcpy(v->data.pData, data, len);
×
979
      qDebug("buildGroupInfo:%d add json data len:%d, data:%s", i, len, (char*)v->data.pData);
×
980
    } else if (IS_VAR_DATA_TYPE(pValue->info.type)) {
651,936✔
981
      if (varDataTLen(data) > pValue->info.bytes) {
448,494✔
982
        code = TSDB_CODE_TDB_INVALID_TABLE_SCHEMA_VER;
×
983
        goto end;
×
984
      }
985
      v->data.type = pValue->info.type;
447,703✔
986
      v->data.nData = varDataLen(data);
447,285✔
987
      v->data.pData = taosMemoryCalloc(1, varDataLen(data) + 1);
447,703✔
988
      if (v->data.pData == NULL) {
447,285✔
989
        code = terrno;
×
990
        goto end;
×
991
      }
992
      memcpy(v->data.pData, varDataVal(data), varDataLen(data));
447,285✔
993
      qDebug("buildGroupInfo:%d add var data type:%d, len:%d, data:%s", i, pValue->info.type, varDataLen(data), (char*)v->data.pData);
447,703✔
994
    } else if (pValue->info.type == TSDB_DATA_TYPE_DECIMAL) {  // reader todo decimal
204,271✔
995
      v->data.type = pValue->info.type;
×
996
      v->data.nData = pValue->info.bytes;
×
997
      v->data.pData = taosMemoryCalloc(1, pValue->info.bytes);
×
998
      if (v->data.pData == NULL) {
×
999
        code = terrno;
×
1000
        goto end;
×
1001
      }
1002
      memcpy(&v->data.pData, data, pValue->info.bytes);
×
1003
      qDebug("buildGroupInfo:%d add data type:%d, data:%"PRId64, i, pValue->info.type, v->data.val);
×
1004
    } else {  // reader todo decimal
1005
      v->data.type = pValue->info.type;
204,671✔
1006
      memcpy(&v->data.val, data, pValue->info.bytes);
205,031✔
1007
      qDebug("buildGroupInfo:%d add data type:%d, data:%"PRId64, i, pValue->info.type, v->data.val);
205,031✔
1008
    }
1009
  }
1010
end:
79,522✔
1011
  if (code != TSDB_CODE_SUCCESS) {
667,134✔
1012
    qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(code));
×
1013
    v->isNull = true;
×
1014
  }
1015
  return code;
667,134✔
1016
}
1017

1018
static void getColInfoResultForGroupbyForStream(void* pVnode, SNodeList* group, STableListInfo* pTableListInfo,
106,340✔
1019
                                   SStorageAPI* pAPI, SHashObj* groupIdMap) {
1020
  int32_t      code = TSDB_CODE_SUCCESS;
106,340✔
1021
  int32_t      lino = 0;
106,340✔
1022
  SArray*      pBlockList = NULL;
106,340✔
1023
  SSDataBlock* pResBlock = NULL;
106,340✔
1024
  SArray*      groupData = NULL;
106,748✔
1025
  SArray*      pUidTagList = NULL;
106,748✔
1026
  SArray*      gInfo = NULL;
106,748✔
1027
  int32_t      tbNameIndex = 0;
106,748✔
1028

1029
  int32_t rows = taosArrayGetSize(pTableListInfo->pTableList);
106,748✔
1030
  if (rows == 0) {
106,748✔
1031
    return;
×
1032
  }
1033

1034
  pUidTagList = taosArrayInit(8, sizeof(STUidTagInfo));
106,748✔
1035
  QUERY_CHECK_NULL(pUidTagList, code, lino, end, terrno);
106,748✔
1036

1037
  for (int32_t i = 0; i < rows; ++i) {
450,268✔
1038
    STableKeyInfo* pkeyInfo = taosArrayGet(pTableListInfo->pTableList, i);
343,520✔
1039
    QUERY_CHECK_NULL(pkeyInfo, code, lino, end, terrno);
343,520✔
1040
    STUidTagInfo info = {.uid = pkeyInfo->uid};
343,520✔
1041
    void*        tmp = taosArrayPush(pUidTagList, &info);
343,520✔
1042
    QUERY_CHECK_NULL(tmp, code, lino, end, terrno);
343,520✔
1043
  }
1044
 
1045
  if (taosArrayGetSize(pUidTagList) > 0) {
106,748✔
1046
    code = pAPI->metaFn.getTableTagsByUid(pVnode, pTableListInfo->idInfo.suid, pUidTagList);
106,748✔
1047
  } else {
1048
    code = pAPI->metaFn.getTableTags(pVnode, pTableListInfo->idInfo.suid, pUidTagList);
×
1049
  }
1050
  if (code != TSDB_CODE_SUCCESS) {
106,748✔
1051
    goto end;
×
1052
  }
1053

1054
  SArray* pColList = NULL;
106,748✔
1055
  code = qGetColumnsFromNodeList(group, true, &pColList);
106,748✔
1056
  if (code != TSDB_CODE_SUCCESS) {
106,748✔
1057
    goto end;
×
1058
  }
1059

1060
  for (int32_t i = 0; i < taosArrayGetSize(pColList); ++i) {
274,442✔
1061
    SColumnInfo* tmp = (SColumnInfo*)taosArrayGet(pColList, i);
167,694✔
1062
    if (tmp != NULL && tmp->colId == -1) {
167,694✔
1063
      tbNameIndex = i;
106,340✔
1064
    }
1065
  }
1066
  
1067
  int32_t numOfTables = taosArrayGetSize(pUidTagList);
106,340✔
1068
  pResBlock = createTagValBlockForFilter(pColList, numOfTables, pUidTagList, pVnode, pAPI);
106,340✔
1069
  taosArrayDestroy(pColList);
106,748✔
1070
  if (pResBlock == NULL) {
106,748✔
1071
    code = terrno;
×
1072
    goto end;
×
1073
  }
1074

1075
  pBlockList = taosArrayInit(2, POINTER_BYTES);
106,748✔
1076
  QUERY_CHECK_NULL(pBlockList, code, lino, end, terrno);
106,748✔
1077

1078
  void* tmp = taosArrayPush(pBlockList, &pResBlock);
106,748✔
1079
  QUERY_CHECK_NULL(tmp, code, lino, end, terrno);
106,748✔
1080

1081
  groupData = taosArrayInit(2, POINTER_BYTES);
106,748✔
1082
  QUERY_CHECK_NULL(groupData, code, lino, end, terrno);
106,748✔
1083

1084
  SNode* pNode = NULL;
106,748✔
1085
  FOREACH(pNode, group) {
274,850✔
1086
    SScalarParam output = {0};
168,102✔
1087

1088
    switch (nodeType(pNode)) {
168,102✔
1089
      case QUERY_NODE_VALUE:
×
1090
        break;
×
1091
      case QUERY_NODE_COLUMN:
168,102✔
1092
      case QUERY_NODE_OPERATOR:
1093
      case QUERY_NODE_FUNCTION: {
1094
        SExprNode* expNode = (SExprNode*)pNode;
168,102✔
1095
        code = createResultData(&expNode->resType, rows, &output);
168,102✔
1096
        if (code != TSDB_CODE_SUCCESS) {
168,102✔
1097
          goto end;
×
1098
        }
1099
        break;
168,102✔
1100
      }
1101

1102
      default:
×
1103
        code = TSDB_CODE_OPS_NOT_SUPPORT;
×
1104
        goto end;
×
1105
    }
1106

1107
    if (nodeType(pNode) == QUERY_NODE_COLUMN) {
168,102✔
1108
      SColumnNode*     pSColumnNode = (SColumnNode*)pNode;
168,102✔
1109
      SColumnInfoData* pColInfo = (SColumnInfoData*)taosArrayGet(pResBlock->pDataBlock, pSColumnNode->slotId);
168,102✔
1110
      QUERY_CHECK_NULL(pColInfo, code, lino, end, terrno);
168,102✔
1111
      code = colDataAssign(output.columnData, pColInfo, rows, NULL);
168,102✔
1112
    } else if (nodeType(pNode) == QUERY_NODE_VALUE) {
×
1113
      continue;
×
1114
    } else {
1115
      gTaskScalarExtra.pStreamInfo = NULL;
×
1116
      gTaskScalarExtra.pStreamRange = NULL;
×
1117
      code = scalarCalculate(pNode, pBlockList, &output, &gTaskScalarExtra);
×
1118
    }
1119

1120
    if (code != TSDB_CODE_SUCCESS) {
168,102✔
1121
      releaseColInfoData(output.columnData);
×
1122
      goto end;
×
1123
    }
1124

1125
    void* tmp = taosArrayPush(groupData, &output.columnData);
168,102✔
1126
    QUERY_CHECK_NULL(tmp, code, lino, end, terrno);
168,102✔
1127
  }
1128

1129
  for (int i = 0; i < rows; i++) {
450,268✔
1130
    gInfo = taosArrayInit(taosArrayGetSize(groupData), sizeof(SStreamGroupValue));
343,520✔
1131
    QUERY_CHECK_NULL(gInfo, code, lino, end, terrno);
343,520✔
1132

1133
    STableKeyInfo* info = taosArrayGet(pTableListInfo->pTableList, i);
343,520✔
1134
    QUERY_CHECK_NULL(info, code, lino, end, terrno);
343,520✔
1135

1136
    for (int j = 0; j < taosArrayGetSize(groupData); j++) {
854,056✔
1137
      SColumnInfoData* pValue = (SColumnInfoData*)taosArrayGetP(groupData, j);
510,536✔
1138
        int32_t ret = buildGroupInfo(pValue, i, gInfo);
510,536✔
1139
        if (ret != TSDB_CODE_SUCCESS) {
510,536✔
1140
          qError("buildGroupInfo failed at line %d since %s", __LINE__, tstrerror(ret));
×
1141
          goto end;
×
1142
        }
1143
        if (j == tbNameIndex) {
510,536✔
1144
          SStreamGroupValue* v = taosArrayGetLast(gInfo);
343,520✔
1145
          if (v != NULL){
343,520✔
1146
            v->isTbname = true;
343,520✔
1147
            v->uid = info->uid;
343,520✔
1148
          }
1149
        }
1150
    }
1151

1152
    int32_t ret = taosHashPut(groupIdMap, &info->uid, sizeof(info->uid), &gInfo, POINTER_BYTES);
343,520✔
1153
    if (ret != TSDB_CODE_SUCCESS) {
343,520✔
1154
      qError("put groupid to map failed at line %d since %s", __LINE__, tstrerror(ret));
×
1155
      goto end;
×
1156
    }
1157
    qDebug("put groupid to map gid:%" PRIu64, info->uid);
343,520✔
1158
    gInfo = NULL;
343,520✔
1159
  }
1160

1161
end:
106,748✔
1162
  blockDataDestroy(pResBlock);
106,748✔
1163
  taosArrayDestroy(pBlockList);
106,748✔
1164
  taosArrayDestroyEx(pUidTagList, freeItem);
106,748✔
1165
  taosArrayDestroyP(groupData, releaseColInfoData);
106,748✔
1166
  taosArrayDestroyEx(gInfo, tDestroySStreamGroupValue);
106,748✔
1167

1168
  if (code != TSDB_CODE_SUCCESS) {
106,748✔
1169
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
1170
  }
1171
}
1172

1173
int32_t getColInfoResultForGroupby(void* pVnode, SNodeList* group, STableListInfo* pTableListInfo, uint8_t* digest,
1,770,317✔
1174
                                   SStorageAPI* pAPI, bool initRemainGroups, SHashObj* groupIdMap) {
1175
  int32_t      code = TSDB_CODE_SUCCESS;
1,770,317✔
1176
  int32_t      lino = 0;
1,770,317✔
1177
  SArray*      pBlockList = NULL;
1,770,317✔
1178
  SSDataBlock* pResBlock = NULL;
1,770,317✔
1179
  void*        keyBuf = NULL;
1,770,935✔
1180
  SArray*      groupData = NULL;
1,770,935✔
1181
  SArray*      pUidTagList = NULL;
1,770,935✔
1182
  SArray*      tableList = NULL;
1,770,935✔
1183
  SArray*      gInfo = NULL;
1,770,935✔
1184

1185
  int32_t rows = taosArrayGetSize(pTableListInfo->pTableList);
1,770,935✔
1186
  if (rows == 0) {
1,770,935✔
1187
    return TSDB_CODE_SUCCESS;
×
1188
  } 
1189

1190
  T_MD5_CTX context = {0};
1,770,935✔
1191
  if (tsTagFilterCache && groupIdMap == NULL) {
1,770,935✔
1192
    SNodeListNode* listNode = NULL;
×
1193
    code = nodesMakeNode(QUERY_NODE_NODE_LIST, (SNode**)&listNode);
×
1194
    if (TSDB_CODE_SUCCESS != code) {
×
1195
      goto end;
×
1196
    }
1197
    listNode->pNodeList = group;
×
1198
    code = genTbGroupDigest((SNode*)listNode, digest, &context);
×
1199
    QUERY_CHECK_CODE(code, lino, end);
×
1200

1201
    nodesFree(listNode);
×
1202

1203
    code = pAPI->metaFn.metaGetCachedTbGroup(pVnode, pTableListInfo->idInfo.suid, context.digest,
×
1204
                                             tListLen(context.digest), &tableList);
1205
    QUERY_CHECK_CODE(code, lino, end);
×
1206

1207
    if (tableList) {
×
1208
      taosArrayDestroy(pTableListInfo->pTableList);
×
1209
      pTableListInfo->pTableList = tableList;
×
1210
      qDebug("retrieve tb group list from cache, numOfTables:%d",
×
1211
             (int32_t)taosArrayGetSize(pTableListInfo->pTableList));
1212
      goto end;
×
1213
    }
1214
  }
1215

1216
  pUidTagList = taosArrayInit(8, sizeof(STUidTagInfo));
1,770,935✔
1217
  QUERY_CHECK_NULL(pUidTagList, code, lino, end, terrno);
1,770,317✔
1218

1219
  for (int32_t i = 0; i < rows; ++i) {
10,463,022✔
1220
    STableKeyInfo* pkeyInfo = taosArrayGet(pTableListInfo->pTableList, i);
8,691,673✔
1221
    QUERY_CHECK_NULL(pkeyInfo, code, lino, end, terrno);
8,691,673✔
1222
    STUidTagInfo info = {.uid = pkeyInfo->uid};
8,691,673✔
1223
    void*        tmp = taosArrayPush(pUidTagList, &info);
8,692,705✔
1224
    QUERY_CHECK_NULL(tmp, code, lino, end, terrno);
8,692,705✔
1225
  }
1226

1227
  if (taosArrayGetSize(pUidTagList) > 0) {
1,771,349✔
1228
    code = pAPI->metaFn.getTableTagsByUid(pVnode, pTableListInfo->idInfo.suid, pUidTagList);
1,770,935✔
1229
  } else {
1230
    code = pAPI->metaFn.getTableTags(pVnode, pTableListInfo->idInfo.suid, pUidTagList);
×
1231
  }
1232
  if (code != TSDB_CODE_SUCCESS) {
1,770,935✔
1233
    goto end;
×
1234
  }
1235

1236
  SArray* pColList = NULL;
1,770,935✔
1237
  code = qGetColumnsFromNodeList(group, true, &pColList); 
1,770,935✔
1238
  if (code != TSDB_CODE_SUCCESS) {
1,768,934✔
1239
    goto end;
×
1240
  }
1241

1242
  int32_t numOfTables = taosArrayGetSize(pUidTagList);
1,768,934✔
1243
  pResBlock = createTagValBlockForFilter(pColList, numOfTables, pUidTagList, pVnode, pAPI);
1,768,934✔
1244
  taosArrayDestroy(pColList);
1,770,317✔
1245
  if (pResBlock == NULL) {
1,770,317✔
1246
    code = terrno;
×
1247
    goto end;
×
1248
  }
1249

1250
  //  int64_t st1 = taosGetTimestampUs();
1251
  //  qDebug("generate tag block rows:%d, cost:%ld us", rows, st1-st);
1252

1253
  pBlockList = taosArrayInit(2, POINTER_BYTES);
1,770,317✔
1254
  QUERY_CHECK_NULL(pBlockList, code, lino, end, terrno);
1,770,317✔
1255

1256
  void* tmp = taosArrayPush(pBlockList, &pResBlock);
1,770,935✔
1257
  QUERY_CHECK_NULL(tmp, code, lino, end, terrno);
1,770,935✔
1258

1259
  groupData = taosArrayInit(2, POINTER_BYTES);
1,770,935✔
1260
  QUERY_CHECK_NULL(groupData, code, lino, end, terrno);
1,769,883✔
1261

1262
  SNode* pNode = NULL;
1,769,883✔
1263
  FOREACH(pNode, group) {
3,661,811✔
1264
    SScalarParam output = {0};
1,893,666✔
1265

1266
    switch (nodeType(pNode)) {
1,893,666✔
1267
      case QUERY_NODE_VALUE:
×
1268
        break;
×
1269
      case QUERY_NODE_COLUMN:
1,885,959✔
1270
      case QUERY_NODE_OPERATOR:
1271
      case QUERY_NODE_FUNCTION: {
1272
        SExprNode* expNode = (SExprNode*)pNode;
1,885,959✔
1273
        code = createResultData(&expNode->resType, rows, &output);
1,885,959✔
1274
        if (code != TSDB_CODE_SUCCESS) {
1,885,557✔
1275
          goto end;
×
1276
        }
1277
        break;
1,885,557✔
1278
      }
1279
      case QUERY_NODE_REMOTE_VALUE: {
7,416✔
1280
        SRemoteValueNode* pRemote = (SRemoteValueNode*)pNode;
7,416✔
1281
        code = qFetchRemoteValue(gTaskScalarExtra.pSubJobCtx, pRemote->subQIdx, pRemote);
7,416✔
1282
        QUERY_CHECK_CODE(code, lino, end);
7,416✔
1283
        break;
7,416✔
1284
      }
1285
      
1286
      default:
×
1287
        code = TSDB_CODE_OPS_NOT_SUPPORT;
×
1288
        goto end;
×
1289
    }
1290

1291
    if (nodeType(pNode) == QUERY_NODE_COLUMN) {
1,892,973✔
1292
      SColumnNode*     pSColumnNode = (SColumnNode*)pNode;
1,869,422✔
1293
      SColumnInfoData* pColInfo = (SColumnInfoData*)taosArrayGet(pResBlock->pDataBlock, pSColumnNode->slotId);
1,869,422✔
1294
      QUERY_CHECK_NULL(pColInfo, code, lino, end, terrno);
1,866,847✔
1295
      code = colDataAssign(output.columnData, pColInfo, rows, NULL);
1,866,847✔
1296
    } else if (nodeType(pNode) == QUERY_NODE_VALUE) {
23,631✔
1297
      continue;
7,416✔
1298
    } else {
1299
      gTaskScalarExtra.pStreamInfo = NULL;
16,971✔
1300
      gTaskScalarExtra.pStreamRange = NULL;
16,971✔
1301
      code = scalarCalculate(pNode, pBlockList, &output, &gTaskScalarExtra);
16,971✔
1302
    }
1303

1304
    if (code != TSDB_CODE_SUCCESS) {
1,884,650✔
1305
      releaseColInfoData(output.columnData);
×
1306
      goto end;
×
1307
    }
1308

1309
    void* tmp = taosArrayPush(groupData, &output.columnData);
1,885,130✔
1310
    QUERY_CHECK_NULL(tmp, code, lino, end, terrno);
1,885,130✔
1311
  }
1312

1313
  int32_t keyLen = 0;
1,769,130✔
1314
  SNode*  node;
1315
  FOREACH(node, group) {
3,662,495✔
1316
    SExprNode* pExpr = (SExprNode*)node;
1,892,023✔
1317
    keyLen += pExpr->resType.bytes;
1,892,023✔
1318
  }
1319

1320
  int32_t nullFlagSize = sizeof(int8_t) * LIST_LENGTH(group);
1,769,071✔
1321
  keyLen += nullFlagSize;
1,768,928✔
1322

1323
  keyBuf = taosMemoryCalloc(1, keyLen);
1,768,928✔
1324
  if (keyBuf == NULL) {
1,768,610✔
1325
    code = terrno;
×
1326
    goto end;
×
1327
  }
1328

1329
  if (initRemainGroups) {
1,768,610✔
1330
    pTableListInfo->remainGroups =
805,975✔
1331
        taosHashInit(rows, taosGetDefaultHashFunction(TSDB_DATA_TYPE_BIGINT), false, HASH_NO_LOCK);
805,699✔
1332
    if (pTableListInfo->remainGroups == NULL) {
805,975✔
1333
      code = terrno;
×
1334
      goto end;
×
1335
    }
1336
  }
1337

1338
  for (int i = 0; i < rows; i++) {
10,462,618✔
1339
    STableKeyInfo* info = taosArrayGet(pTableListInfo->pTableList, i);
8,692,103✔
1340
    QUERY_CHECK_NULL(info, code, lino, end, terrno);
8,690,505✔
1341

1342
    if (groupIdMap != NULL){
8,690,505✔
1343
      gInfo = taosArrayInit(taosArrayGetSize(groupData), sizeof(SStreamGroupValue));
132,937✔
1344
    }
1345
    
1346
    char* isNull = (char*)keyBuf;
8,691,440✔
1347
    char* pStart = (char*)keyBuf + sizeof(int8_t) * LIST_LENGTH(group);
8,691,440✔
1348
    for (int j = 0; j < taosArrayGetSize(groupData); j++) {
18,018,254✔
1349
      SColumnInfoData* pValue = (SColumnInfoData*)taosArrayGetP(groupData, j);
9,329,691✔
1350

1351
      if (groupIdMap != NULL && gInfo != NULL) {
9,329,410✔
1352
        int32_t ret = buildGroupInfo(pValue, i, gInfo);
156,598✔
1353
        if (ret != TSDB_CODE_SUCCESS) {
156,598✔
1354
          qError("buildGroupInfo failed at line %d since %s", __LINE__, tstrerror(ret));
×
1355
          taosArrayDestroyEx(gInfo, tDestroySStreamGroupValue);
×
1356
          gInfo = NULL;
×
1357
        }
1358
      }
1359
      
1360
      if (colDataIsNull_s(pValue, i)) {
18,660,703✔
1361
        isNull[j] = 1;
94,750✔
1362
      } else {
1363
        isNull[j] = 0;
9,236,543✔
1364
        char* data = colDataGetData(pValue, i);
9,233,829✔
1365
        if (pValue->info.type == TSDB_DATA_TYPE_JSON) {
9,233,373✔
1366
          // if (tTagIsJson(data)) {
1367
          //   code = TSDB_CODE_QRY_JSON_IN_GROUP_ERROR;
1368
          //   goto end;
1369
          // }
1370
          if (tTagIsJsonNull(data)) {
89,813✔
1371
            isNull[j] = 1;
×
1372
            continue;
×
1373
          }
1374
          int32_t len = getJsonValueLen(data);
89,813✔
1375
          memcpy(pStart, data, len);
89,813✔
1376
          pStart += len;
89,813✔
1377
        } else if (IS_VAR_DATA_TYPE(pValue->info.type)) {
9,140,722✔
1378
          if (IS_STR_DATA_BLOB(pValue->info.type)) {
6,134,434✔
1379
            if (blobDataTLen(data) > TSDB_MAX_BLOB_LEN) {
710✔
1380
              code = TSDB_CODE_TDB_INVALID_TABLE_SCHEMA_VER;
×
1381
              goto end;
×
1382
            }
1383
            memcpy(pStart, data, blobDataTLen(data));
×
1384
            pStart += blobDataTLen(data);
×
1385
          } else {
1386
            if (varDataTLen(data) > pValue->info.bytes) {
6,129,049✔
1387
              code = TSDB_CODE_TDB_INVALID_TABLE_SCHEMA_VER;
×
1388
              goto end;
×
1389
            }
1390
            memcpy(pStart, data, varDataTLen(data));
6,129,748✔
1391
            pStart += varDataTLen(data);
6,130,310✔
1392
          }
1393
        } else {
1394
          memcpy(pStart, data, pValue->info.bytes);
3,013,119✔
1395
          pStart += pValue->info.bytes;
3,015,694✔
1396
        }
1397
      }
1398
    }
1399

1400
    int32_t len = (int32_t)(pStart - (char*)keyBuf);
8,686,055✔
1401
    info->groupId = calcGroupId(keyBuf, len);
8,686,055✔
1402
    if (groupIdMap != NULL && gInfo != NULL) {
8,683,213✔
1403
      int32_t ret = taosHashPut(groupIdMap, &info->groupId, sizeof(info->groupId), &gInfo, POINTER_BYTES);
132,550✔
1404
      if (ret != TSDB_CODE_SUCCESS) {
132,937✔
1405
        qError("put groupid to map failed at line %d since %s", __LINE__, tstrerror(ret));
×
1406
        taosArrayDestroyEx(gInfo, tDestroySStreamGroupValue);
×
1407
      }
1408
      qDebug("put groupid to map gid:%" PRIu64, info->groupId);
132,937✔
1409
      gInfo = NULL;
132,937✔
1410
    }
1411
    if (initRemainGroups) {
8,683,600✔
1412
      // groupId ~ table uid
1413
      code = taosHashPut(pTableListInfo->remainGroups, &(info->groupId), sizeof(info->groupId), &(info->uid),
4,275,625✔
1414
                         sizeof(info->uid));
1415
      if (code == TSDB_CODE_DUP_KEY) {
4,284,068✔
1416
        code = TSDB_CODE_SUCCESS;
848,257✔
1417
      }
1418
      QUERY_CHECK_CODE(code, lino, end);
4,284,068✔
1419
    }
1420
  }
1421

1422
  if (tsTagFilterCache && groupIdMap == NULL) {
1,770,515✔
1423
    tableList = taosArrayDup(pTableListInfo->pTableList, NULL);
×
1424
    QUERY_CHECK_NULL(tableList, code, lino, end, terrno);
×
1425

1426
    code = pAPI->metaFn.metaPutTbGroupToCache(pVnode, pTableListInfo->idInfo.suid, context.digest,
×
1427
                                              tListLen(context.digest), tableList,
1428
                                              taosArrayGetSize(tableList) * sizeof(STableKeyInfo));
×
1429
    QUERY_CHECK_CODE(code, lino, end);
×
1430
  }
1431

1432
  //  int64_t st2 = taosGetTimestampUs();
1433
  //  qDebug("calculate tag block rows:%d, cost:%ld us", rows, st2-st1);
1434

1435
end:
1,769,328✔
1436
  taosMemoryFreeClear(keyBuf);
1,770,935✔
1437
  blockDataDestroy(pResBlock);
1,770,935✔
1438
  taosArrayDestroy(pBlockList);
1,770,511✔
1439
  taosArrayDestroyEx(pUidTagList, freeItem);
1,770,649✔
1440
  taosArrayDestroyP(groupData, releaseColInfoData);
1,770,792✔
1441
  taosArrayDestroyEx(gInfo, tDestroySStreamGroupValue);
1,770,792✔
1442

1443
  if (code != TSDB_CODE_SUCCESS) {
1,770,792✔
1444
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
1445
  }
1446
  return code;
1,770,654✔
1447
}
1448

1449
static int32_t nameComparFn(const void* p1, const void* p2) {
928,208✔
1450
  const char* pName1 = *(const char**)p1;
928,208✔
1451
  const char* pName2 = *(const char**)p2;
928,208✔
1452

1453
  int32_t ret = strcmp(pName1, pName2);
928,208✔
1454
  if (ret == 0) {
928,208✔
1455
    return 0;
18,372✔
1456
  } else {
1457
    return (ret > 0) ? 1 : -1;
909,836✔
1458
  }
1459
}
1460

1461
static SArray* getTableNameList(const SNodeListNode* pList) {
526,675✔
1462
  int32_t    code = TSDB_CODE_SUCCESS;
526,675✔
1463
  int32_t    lino = 0;
526,675✔
1464
  int32_t    len = LIST_LENGTH(pList->pNodeList);
526,675✔
1465
  SListCell* cell = pList->pNodeList->pHead;
526,675✔
1466

1467
  SArray* pTbList = taosArrayInit(len, POINTER_BYTES);
526,125✔
1468
  QUERY_CHECK_NULL(pTbList, code, lino, _end, terrno);
526,675✔
1469

1470
  for (int i = 0; i < pList->pNodeList->length; i++) {
1,458,129✔
1471
    SValueNode* valueNode = (SValueNode*)cell->pNode;
930,904✔
1472
    if (!IS_VAR_DATA_TYPE(valueNode->node.resType.type)) {
930,904✔
1473
      terrno = TSDB_CODE_INVALID_PARA;
×
1474
      taosArrayDestroy(pTbList);
×
1475
      return NULL;
×
1476
    }
1477

1478
    char* name = varDataVal(valueNode->datum.p);
931,454✔
1479
    void* tmp = taosArrayPush(pTbList, &name);
931,454✔
1480
    QUERY_CHECK_NULL(tmp, code, lino, _end, terrno);
931,454✔
1481
    cell = cell->pNext;
931,454✔
1482
  }
1483

1484
  size_t numOfTables = taosArrayGetSize(pTbList);
526,675✔
1485

1486
  // order the name
1487
  taosArraySort(pTbList, nameComparFn);
526,675✔
1488

1489
  // remove the duplicates
1490
  SArray* pNewList = taosArrayInit(taosArrayGetSize(pTbList), sizeof(void*));
526,675✔
1491
  QUERY_CHECK_NULL(pNewList, code, lino, _end, terrno);
526,675✔
1492
  void* tmpTbl = taosArrayGet(pTbList, 0);
526,675✔
1493
  QUERY_CHECK_NULL(tmpTbl, code, lino, _end, terrno);
526,675✔
1494
  void* tmp = taosArrayPush(pNewList, tmpTbl);
526,675✔
1495
  QUERY_CHECK_NULL(tmp, code, lino, _end, terrno);
526,675✔
1496

1497
  for (int32_t i = 1; i < numOfTables; ++i) {
931,454✔
1498
    char** name = taosArrayGetLast(pNewList);
404,779✔
1499
    char** nameInOldList = taosArrayGet(pTbList, i);
404,779✔
1500
    QUERY_CHECK_NULL(nameInOldList, code, lino, _end, terrno);
404,779✔
1501
    if (strcmp(*name, *nameInOldList) == 0) {
404,779✔
1502
      continue;
9,884✔
1503
    }
1504

1505
    tmp = taosArrayPush(pNewList, nameInOldList);
394,895✔
1506
    QUERY_CHECK_NULL(tmp, code, lino, _end, terrno);
394,895✔
1507
  }
1508

1509
_end:
526,675✔
1510
  taosArrayDestroy(pTbList);
526,675✔
1511
  if (code != TSDB_CODE_SUCCESS) {
526,675✔
1512
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
1513
    return NULL;
×
1514
  }
1515
  return pNewList;
526,675✔
1516
}
1517

1518
static int tableUidCompare(const void* a, const void* b) {
×
1519
  uint64_t u1 = *(uint64_t*)a;
×
1520
  uint64_t u2 = *(uint64_t*)b;
×
1521

1522
  if (u1 == u2) {
×
1523
    return 0;
×
1524
  }
1525

1526
  return u1 < u2 ? -1 : 1;
×
1527
}
1528

1529
static int32_t filterTableInfoCompare(const void* a, const void* b) {
18,473,235✔
1530
  STUidTagInfo* p1 = (STUidTagInfo*)a;
18,473,235✔
1531
  STUidTagInfo* p2 = (STUidTagInfo*)b;
18,473,235✔
1532

1533
  if (p1->uid == p2->uid) {
18,473,235✔
1534
    return 0;
×
1535
  }
1536

1537
  return p1->uid < p2->uid ? -1 : 1;
18,473,235✔
1538
}
1539

1540
static FilterCondType checkTagCond(SNode* cond) {
12,668,282✔
1541
  if (nodeType(cond) == QUERY_NODE_OPERATOR) {
12,668,282✔
1542
    return FILTER_NO_LOGIC;
10,728,920✔
1543
  }
1544
  if (nodeType(cond) == QUERY_NODE_LOGIC_CONDITION && ((SLogicConditionNode*)cond)->condType == LOGIC_COND_TYPE_AND) {
1,939,362✔
1545
    return FILTER_AND;
1,720,945✔
1546
  }
1547
  return FILTER_OTHER;
217,417✔
1548
}
1549

1550
static int32_t optimizeTbnameInCond(void* pVnode, int64_t suid, SArray* list, SNode* cond, SStorageAPI* pAPI) {
13,137,787✔
1551
  int32_t ret = -1;
13,137,787✔
1552
  int32_t ntype = nodeType(cond);
13,137,787✔
1553

1554
  if (ntype == QUERY_NODE_OPERATOR) {
13,137,787✔
1555
    ret = optimizeTbnameInCondImpl(pVnode, list, cond, pAPI, suid);
11,192,132✔
1556
    return ret;
11,191,403✔
1557
  }
1558
  if (ntype != QUERY_NODE_LOGIC_CONDITION || ((SLogicConditionNode*)cond)->condType != LOGIC_COND_TYPE_AND) {
1,945,655✔
1559
    return ret;
217,641✔
1560
  }
1561

1562
  bool                 hasTbnameCond = false;
1,727,626✔
1563
  SLogicConditionNode* pNode = (SLogicConditionNode*)cond;
1,727,626✔
1564
  SNodeList*           pList = (SNodeList*)pNode->pParameterList;
1,727,626✔
1565

1566
  int32_t len = LIST_LENGTH(pList);
1,728,226✔
1567
  if (len <= 0) {
1,728,032✔
1568
    return ret;
×
1569
  }
1570

1571
  SListCell* cell = pList->pHead;
1,728,032✔
1572
  for (int i = 0; i < len; i++) {
5,609,529✔
1573
    if (cell == NULL) break;
3,888,394✔
1574
    if (optimizeTbnameInCondImpl(pVnode, list, cell->pNode, pAPI, suid) == 0) {
3,888,394✔
1575
      hasTbnameCond = true;
6,293✔
1576
      break;
6,293✔
1577
    }
1578
    cell = cell->pNext;
3,881,493✔
1579
  }
1580

1581
  taosArraySort(list, filterTableInfoCompare);
1,727,428✔
1582
  taosArrayRemoveDuplicate(list, filterTableInfoCompare, NULL);
1,727,820✔
1583

1584
  if (hasTbnameCond) {
1,727,823✔
1585
    ret = pAPI->metaFn.getTableTagsByUid(pVnode, suid, list);
6,293✔
1586
  }
1587

1588
  return ret;
1,727,823✔
1589
}
1590

1591
// only return uid that does not contained in pExistedUidList
1592
static int32_t optimizeTbnameInCondImpl(void* pVnode, SArray* pExistedUidList, SNode* pTagCond, SStorageAPI* pStoreAPI,
15,078,884✔
1593
                                        uint64_t suid) {
1594
  if (nodeType(pTagCond) != QUERY_NODE_OPERATOR) {
15,078,884✔
1595
    return -1;
6,418✔
1596
  }
1597

1598
  SOperatorNode* pNode = (SOperatorNode*)pTagCond;
15,073,063✔
1599
  if (pNode->opType != OP_TYPE_IN) {
15,073,063✔
1600
    return -1;
14,076,526✔
1601
  }
1602

1603
  if ((pNode->pLeft != NULL && ((nodeType(pNode->pLeft) == QUERY_NODE_FUNCTION &&
996,898✔
1604
                                 ((SFunctionNode*)pNode->pLeft)->funcType == FUNCTION_TYPE_TBNAME)) ||
526,675✔
1605
       (nodeType(pNode->pLeft) == QUERY_NODE_COLUMN && ((SColumnNode*)pNode->pLeft)->colType == COLUMN_TYPE_TBNAME)) &&
470,773✔
1606
      (pNode->pRight != NULL && nodeType(pNode->pRight) == QUERY_NODE_NODE_LIST)) {
526,478✔
1607
    SNodeListNode* pList = (SNodeListNode*)pNode->pRight;
526,675✔
1608

1609
    int32_t len = LIST_LENGTH(pList->pNodeList);
526,675✔
1610
    if (len <= 0) {
526,675✔
1611
      return -1;
×
1612
    }
1613

1614
    SArray*   pTbList = getTableNameList(pList);
526,675✔
1615
    int32_t   numOfTables = taosArrayGetSize(pTbList);
526,675✔
1616
    SHashObj* uHash = NULL;
526,675✔
1617

1618
    size_t numOfExisted = taosArrayGetSize(pExistedUidList);  // len > 0 means there already have uids
526,675✔
1619
    if (numOfExisted > 0) {
526,675✔
1620
      uHash = taosHashInit(numOfExisted / 0.7, taosGetDefaultHashFunction(TSDB_DATA_TYPE_BIGINT), false, HASH_NO_LOCK);
2,364✔
1621
      if (!uHash) {
2,364✔
1622
        qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
1623
        return terrno;
×
1624
      }
1625

1626
      for (int i = 0; i < numOfExisted; i++) {
2,361,640✔
1627
        STUidTagInfo* pTInfo = taosArrayGet(pExistedUidList, i);
2,359,276✔
1628
        if (!pTInfo) {
2,359,276✔
1629
          qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
1630
          return terrno;
×
1631
        }
1632
        int32_t tempRes = taosHashPut(uHash, &pTInfo->uid, sizeof(uint64_t), &i, sizeof(i));
2,359,276✔
1633
        if (tempRes != TSDB_CODE_SUCCESS && tempRes != TSDB_CODE_DUP_KEY) {
2,359,276✔
1634
          qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(tempRes));
×
1635
          return tempRes;
×
1636
        }
1637
      }
1638
    }
1639

1640
    for (int i = 0; i < numOfTables; i++) {
1,330,069✔
1641
      char* name = taosArrayGetP(pTbList, i);
860,564✔
1642

1643
      uint64_t uid = 0, csuid = 0;
861,086✔
1644
      if (pStoreAPI->metaFn.getTableUidByName(pVnode, name, &uid) == 0) {
861,086✔
1645
        ETableType tbType = TSDB_TABLE_MAX;
476,184✔
1646
        if (pStoreAPI->metaFn.getTableTypeSuidByName(pVnode, name, &tbType, &csuid) == 0 &&
476,184✔
1647
            tbType == TSDB_CHILD_TABLE) {
475,662✔
1648
          if (suid != csuid) {
418,492✔
1649
            continue;
892✔
1650
          }
1651
          if (NULL == uHash || taosHashGet(uHash, &uid, sizeof(uid)) == NULL) {
417,600✔
1652
            STUidTagInfo s = {.uid = uid, .name = name, .pTagVal = NULL};
416,418✔
1653
            void*        tmp = taosArrayPush(pExistedUidList, &s);
416,940✔
1654
            if (!tmp) {
416,940✔
1655
              return terrno;
×
1656
            }
1657
          }
1658
        } else {
1659
          taosArrayDestroy(pTbList);
57,170✔
1660
          taosHashCleanup(uHash);
57,170✔
1661
          return -1;
57,170✔
1662
        }
1663
      } else {
1664
        //        qWarn("failed to get tableIds from by table name: %s, reason: %s", name, tstrerror(terrno));
1665
        terrno = 0;
384,902✔
1666
      }
1667
    }
1668

1669
    taosHashCleanup(uHash);
469,505✔
1670
    taosArrayDestroy(pTbList);
469,505✔
1671
    return 0;
469,505✔
1672
  }
1673

1674
  return -1;
471,179✔
1675
}
1676

1677
SSDataBlock* createTagValBlockForFilter(SArray* pColList, int32_t numOfTables, SArray* pUidTagList, void* pVnode,
13,667,365✔
1678
                                        SStorageAPI* pStorageAPI) {
1679
  int32_t      code = TSDB_CODE_SUCCESS;
13,667,365✔
1680
  int32_t      lino = 0;
13,667,365✔
1681
  SSDataBlock* pResBlock = NULL;
13,667,365✔
1682
  code = createDataBlock(&pResBlock);
13,669,996✔
1683
  QUERY_CHECK_CODE(code, lino, _end);
13,668,831✔
1684

1685
  for (int32_t i = 0; i < taosArrayGetSize(pColList); ++i) {
28,489,269✔
1686
    SColumnInfoData colInfo = {0};
14,818,099✔
1687
    void*           tmp = taosArrayGet(pColList, i);
14,813,240✔
1688
    QUERY_CHECK_NULL(tmp, code, lino, _end, terrno);
14,809,415✔
1689
    colInfo.info = *(SColumnInfo*)tmp;
14,809,415✔
1690
    code = blockDataAppendColInfo(pResBlock, &colInfo);
14,805,532✔
1691
    QUERY_CHECK_CODE(code, lino, _end);
14,817,522✔
1692
  }
1693

1694
  code = blockDataEnsureCapacity(pResBlock, numOfTables);
13,664,745✔
1695
  if (code != TSDB_CODE_SUCCESS) {
13,664,699✔
1696
    terrno = code;
×
1697
    blockDataDestroy(pResBlock);
×
1698
    return NULL;
×
1699
  }
1700

1701
  pResBlock->info.rows = numOfTables;
13,664,699✔
1702

1703
  int32_t numOfCols = taosArrayGetSize(pResBlock->pDataBlock);
13,663,560✔
1704

1705
  for (int32_t i = 0; i < numOfTables; i++) {
203,631,847✔
1706
    STUidTagInfo* p1 = taosArrayGet(pUidTagList, i);
189,959,487✔
1707
    QUERY_CHECK_NULL(p1, code, lino, _end, terrno);
189,969,506✔
1708

1709
    for (int32_t j = 0; j < numOfCols; j++) {
387,462,402✔
1710
      SColumnInfoData* pColInfo = (SColumnInfoData*)taosArrayGet(pResBlock->pDataBlock, j);
197,450,250✔
1711
      QUERY_CHECK_NULL(pColInfo, code, lino, _end, terrno);
197,448,645✔
1712

1713
      if (pColInfo->info.colId == -1) {  // tbname
197,448,645✔
1714
        char str[TSDB_TABLE_FNAME_LEN + VARSTR_HEADER_SIZE] = {0};
8,186,065✔
1715
        if (p1->name != NULL) {
8,187,381✔
1716
          STR_TO_VARSTR(str, p1->name);
416,422✔
1717
        } else {  // name is not retrieved during filter
1718
          code = pStorageAPI->metaFn.getTableNameByUid(pVnode, p1->uid, str);
7,773,394✔
1719
          QUERY_CHECK_CODE(code, lino, _end);
7,772,620✔
1720
        }
1721

1722
        code = colDataSetVal(pColInfo, i, str, false);
8,189,560✔
1723
        QUERY_CHECK_CODE(code, lino, _end);
8,187,765✔
1724
#if TAG_FILTER_DEBUG
1725
        qDebug("tagfilter uid:%ld, tbname:%s", *uid, str + 2);
1726
#endif
1727
      } else {
1728
        STagVal tagVal = {0};
189,256,207✔
1729
        tagVal.cid = pColInfo->info.colId;
189,272,992✔
1730
        if (p1->pTagVal == NULL) {
189,268,906✔
1731
          colDataSetNULL(pColInfo, i);
9,000✔
1732
        } else {
1733
          const char* p = pStorageAPI->metaFn.extractTagVal(p1->pTagVal, pColInfo->info.type, &tagVal);
189,260,160✔
1734

1735
          if (p == NULL || (pColInfo->info.type == TSDB_DATA_TYPE_JSON && ((STag*)p)->nTag == 0)) {
189,290,043✔
1736
            colDataSetNULL(pColInfo, i);
3,984,797✔
1737
          } else if (pColInfo->info.type == TSDB_DATA_TYPE_JSON) {
185,303,561✔
1738
            code = colDataSetVal(pColInfo, i, p, false);
710,372✔
1739
            QUERY_CHECK_CODE(code, lino, _end);
710,372✔
1740
          } else if (IS_VAR_DATA_TYPE(pColInfo->info.type)) {
299,525,527✔
1741
            if (IS_STR_DATA_BLOB(pColInfo->info.type)) {
114,939,698✔
1742
              QUERY_CHECK_CODE(code = TSDB_CODE_BLOB_NOT_SUPPORT_TAG, lino, _end);
×
1743
            }
1744
            char* tmp = taosMemoryMalloc(tagVal.nData + VARSTR_HEADER_SIZE + 1);
114,936,419✔
1745
            QUERY_CHECK_NULL(tmp, code, lino, _end, terrno);
114,919,497✔
1746
            varDataSetLen(tmp, tagVal.nData);
114,919,497✔
1747
            memcpy(tmp + VARSTR_HEADER_SIZE, tagVal.pData, tagVal.nData);
114,919,113✔
1748
            code = colDataSetVal(pColInfo, i, tmp, false);
114,915,609✔
1749
#if TAG_FILTER_DEBUG
1750
            qDebug("tagfilter varch:%s", tmp + 2);
1751
#endif
1752
            taosMemoryFree(tmp);
114,938,531✔
1753
            QUERY_CHECK_CODE(code, lino, _end);
114,929,757✔
1754
          } else {
1755
            code = colDataSetVal(pColInfo, i, (const char*)&tagVal.i64, false);
69,659,801✔
1756
            QUERY_CHECK_CODE(code, lino, _end);
69,683,490✔
1757
#if TAG_FILTER_DEBUG
1758
            if (pColInfo->info.type == TSDB_DATA_TYPE_INT) {
1759
              qDebug("tagfilter int:%d", *(int*)(&tagVal.i64));
1760
            } else if (pColInfo->info.type == TSDB_DATA_TYPE_DOUBLE) {
1761
              qDebug("tagfilter double:%f", *(double*)(&tagVal.i64));
1762
            }
1763
#endif
1764
          }
1765
        }
1766
      }
1767
    }
1768
  }
1769

1770
_end:
13,670,581✔
1771
  if (code != TSDB_CODE_SUCCESS) {
13,672,360✔
1772
    blockDataDestroy(pResBlock);
1,582✔
1773
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
1774
    terrno = code;
×
1775
    return NULL;
×
1776
  }
1777
  return pResBlock;
13,670,778✔
1778
}
1779

1780
static int32_t doSetQualifiedUid(STableListInfo* pListInfo, SArray* pUidList, const SArray* pUidTagList,
11,134,288✔
1781
                                 bool* pResultList, bool addUid) {
1782
  taosArrayClear(pUidList);
11,134,288✔
1783

1784
  STableKeyInfo info = {.uid = 0, .groupId = 0};
11,131,273✔
1785
  int32_t       numOfTables = taosArrayGetSize(pUidTagList);
11,133,974✔
1786
  for (int32_t i = 0; i < numOfTables; ++i) {
191,105,237✔
1787
    if (pResultList[i]) {
179,955,671✔
1788
      STUidTagInfo* tmpTag = (STUidTagInfo*)taosArrayGet(pUidTagList, i);
78,841,843✔
1789
      if (!tmpTag) {
78,846,426✔
1790
        qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
1791
        return terrno;
×
1792
      }
1793
      uint64_t uid = tmpTag->uid;
78,846,426✔
1794
      qDebug("tagfilter get uid:%" PRId64 ", res:%d", uid, pResultList[i]);
78,844,476✔
1795

1796
      info.uid = uid;
78,864,788✔
1797
      //qInfo("doSetQualifiedUid row:%d added to pTableList", i);
1798
      void* p = taosArrayPush(pListInfo->pTableList, &info);
78,864,788✔
1799
      if (p == NULL) {
78,853,478✔
1800
        return terrno;
×
1801
      }
1802

1803
      if (addUid) {
78,853,478✔
1804
        //qInfo("doSetQualifiedUid row:%d added to pUidList", i);
1805
        void* tmp = taosArrayPush(pUidList, &uid);
19,895✔
1806
        if (tmp == NULL) {
19,895✔
1807
          return terrno;
×
1808
        }
1809
      }
1810
    } else {
1811
      //qInfo("doSetQualifiedUid row:%d failed", i);
1812
    }
1813
  }
1814

1815
  return TSDB_CODE_SUCCESS;
11,149,566✔
1816
}
1817

1818
static int32_t copyExistedUids(SArray* pUidTagList, const SArray* pUidList) {
13,137,118✔
1819
  int32_t code = TSDB_CODE_SUCCESS;
13,137,118✔
1820
  int32_t numOfExisted = taosArrayGetSize(pUidList);
13,137,118✔
1821
  if (numOfExisted == 0) {
13,137,972✔
1822
    return code;
10,235,148✔
1823
  }
1824

1825
  for (int32_t i = 0; i < numOfExisted; ++i) {
35,693,217✔
1826
    uint64_t* uid = taosArrayGet(pUidList, i);
32,789,929✔
1827
    if (!uid) {
32,789,929✔
1828
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
1829
      return terrno;
×
1830
    }
1831
    STUidTagInfo info = {.uid = *uid};
32,789,929✔
1832
    void*        tmp = taosArrayPush(pUidTagList, &info);
32,790,393✔
1833
    if (!tmp) {
32,790,393✔
1834
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
1835
      return code;
×
1836
    }
1837
  }
1838
  return code;
2,903,288✔
1839
}
1840

1841
int32_t doFilterByTagCond(STableListInfo* pListInfo, SArray* pUidList, SNode* pTagCond, void* pVnode,
219,789,745✔
1842
                                 SIdxFltStatus status, SStorageAPI* pAPI, bool addUid, bool* listAdded, void* pStreamInfo) {
1843
  *listAdded = false;
219,789,745✔
1844
  if (pTagCond == NULL) {
219,835,687✔
1845
    return TSDB_CODE_SUCCESS;
206,630,084✔
1846
  }
1847

1848
  terrno = TSDB_CODE_SUCCESS;
13,205,603✔
1849

1850
  int32_t      lino = 0;
13,137,381✔
1851
  int32_t      code = TSDB_CODE_SUCCESS;
13,137,381✔
1852
  SArray*      pBlockList = NULL;
13,137,381✔
1853
  SSDataBlock* pResBlock = NULL;
13,137,381✔
1854
  SScalarParam output = {0};
13,136,175✔
1855
  SArray*      pUidTagList = NULL;
13,136,016✔
1856

1857
  SDataType type = {.type = TSDB_DATA_TYPE_BOOL, .bytes = sizeof(bool)};
13,136,016✔
1858

1859
  //  int64_t stt = taosGetTimestampUs();
1860
  pUidTagList = taosArrayInit(10, sizeof(STUidTagInfo));
13,136,359✔
1861
  QUERY_CHECK_NULL(pUidTagList, code, lino, end, terrno);
13,137,080✔
1862

1863
  code = copyExistedUids(pUidTagList, pUidList);
13,137,080✔
1864
  QUERY_CHECK_CODE(code, lino, end);
13,137,787✔
1865

1866
  int32_t filter = optimizeTbnameInCond(pVnode, pListInfo->idInfo.suid, pUidTagList, pTagCond, pAPI);
13,137,787✔
1867
  if (filter == 0) {  // tbname in filter is activated, do nothing and return
13,136,262✔
1868
    taosArrayClear(pUidList);
469,505✔
1869

1870
    int32_t numOfRows = taosArrayGetSize(pUidTagList);
469,505✔
1871
    code = taosArrayEnsureCap(pUidList, numOfRows);
469,505✔
1872
    QUERY_CHECK_CODE(code, lino, end);
469,505✔
1873

1874
    for (int32_t i = 0; i < numOfRows; ++i) {
3,246,899✔
1875
      STUidTagInfo* pInfo = taosArrayGet(pUidTagList, i);
2,777,394✔
1876
      QUERY_CHECK_NULL(pInfo, code, lino, end, terrno);
2,777,394✔
1877
      void* tmp = taosArrayPush(pUidList, &pInfo->uid);
2,777,394✔
1878
      QUERY_CHECK_NULL(tmp, code, lino, end, terrno);
2,777,394✔
1879
    }
1880
    terrno = 0;
469,505✔
1881
  } else {
1882
    qDebug("pUidTagList size:%d", (int32_t)taosArrayGetSize(pUidTagList));
12,666,757✔
1883

1884
    FilterCondType condType = checkTagCond(pTagCond);
12,666,757✔
1885
    if (((condType == FILTER_NO_LOGIC || condType == FILTER_AND) && status != SFLT_NOT_INDEX) ||
21,924,957✔
1886
          taosArrayGetSize(pUidTagList) > 0) {
9,257,078✔
1887
      code = pAPI->metaFn.getTableTagsByUid(pVnode, pListInfo->idInfo.suid, pUidTagList);
3,695,041✔
1888
    } else {
1889
      code = pAPI->metaFn.getTableTags(pVnode, pListInfo->idInfo.suid, pUidTagList);
8,972,838✔
1890
    }
1891
    if (code != TSDB_CODE_SUCCESS) {
12,665,280✔
1892
      qError("failed to get table tags from meta, reason:%s, suid:%" PRIu64, tstrerror(code), pListInfo->idInfo.suid);
×
1893
      terrno = code;
×
1894
      QUERY_CHECK_CODE(code, lino, end);
×
1895
    }
1896
  }
1897

1898
  qDebug("final pUidTagList size:%d", (int32_t)taosArrayGetSize(pUidTagList));
13,134,785✔
1899

1900
  int32_t numOfTables = taosArrayGetSize(pUidTagList);
13,136,455✔
1901
  if (numOfTables == 0) {
13,137,775✔
1902
    goto end;
1,985,399✔
1903
  }
1904

1905
  SArray* pColList = NULL;
11,152,376✔
1906
  code = qGetColumnsFromNodeList(pTagCond, false, &pColList); 
11,152,376✔
1907
  if (code != TSDB_CODE_SUCCESS) {
11,143,788✔
1908
    goto end;
×
1909
  }
1910
  pResBlock = createTagValBlockForFilter(pColList, numOfTables, pUidTagList, pVnode, pAPI);
11,143,788✔
1911
  taosArrayDestroy(pColList);
11,150,456✔
1912
  if (pResBlock == NULL) {
11,150,914✔
1913
    code = terrno;
×
1914
    QUERY_CHECK_CODE(code, lino, end);
×
1915
  }
1916

1917
  //fprintDataBlock(pResBlock, "tagFilter", "", 0);
1918

1919
  //  int64_t st1 = taosGetTimestampUs();
1920
  //  qDebug("generate tag block rows:%d, cost:%ld us", rows, st1-st);
1921
  pBlockList = taosArrayInit(2, POINTER_BYTES);
11,150,914✔
1922
  QUERY_CHECK_NULL(pBlockList, code, lino, end, terrno);
11,150,135✔
1923

1924
  void* tmp = taosArrayPush(pBlockList, &pResBlock);
11,152,175✔
1925
  QUERY_CHECK_NULL(tmp, code, lino, end, terrno);
11,152,175✔
1926

1927
  code = createResultData(&type, numOfTables, &output);
11,152,175✔
1928
  if (code != TSDB_CODE_SUCCESS) {
11,148,083✔
1929
    terrno = code;
×
1930
    QUERY_CHECK_CODE(code, lino, end);
×
1931
  }
1932

1933
  gTaskScalarExtra.pStreamInfo = pStreamInfo;
11,148,083✔
1934
  gTaskScalarExtra.pStreamRange = NULL;
11,148,083✔
1935
  code = scalarCalculate(pTagCond, pBlockList, &output, &gTaskScalarExtra);
11,146,864✔
1936
  if (code != TSDB_CODE_SUCCESS) {
11,135,509✔
1937
    qError("failed to calculate scalar, reason:%s", tstrerror(code));
1,102✔
1938
    terrno = code;
1,102✔
1939
    QUERY_CHECK_CODE(code, lino, end);
1,102✔
1940
  }
1941

1942
  code = doSetQualifiedUid(pListInfo, pUidList, pUidTagList, (bool*)output.columnData->pData, addUid);
11,134,407✔
1943
  if (code != TSDB_CODE_SUCCESS) {
11,149,925✔
1944
    terrno = code;
×
1945
    QUERY_CHECK_CODE(code, lino, end);
×
1946
  }
1947
  *listAdded = true;
11,149,925✔
1948

1949
end:
13,135,647✔
1950
  if (code != TSDB_CODE_SUCCESS) {
13,133,653✔
1951
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
1,102✔
1952
  }
1953
  blockDataDestroy(pResBlock);
13,133,653✔
1954
  taosArrayDestroy(pBlockList);
13,133,192✔
1955
  taosArrayDestroyEx(pUidTagList, freeItem);
13,129,454✔
1956

1957
  colDataDestroy(output.columnData);
13,135,625✔
1958
  taosMemoryFreeClear(output.columnData);
13,136,306✔
1959
  return code;
13,136,860✔
1960
}
1961

1962
typedef struct {
1963
  int32_t code;
1964
  SStreamRuntimeFuncInfo* pStreamRuntimeInfo;
1965
} PlaceHolderContext;
1966

1967
static EDealRes replacePlaceHolderColumn(SNode** pNode, void* pContext) {
164,292✔
1968
  PlaceHolderContext* pData = (PlaceHolderContext*)pContext;
164,292✔
1969
  if (QUERY_NODE_FUNCTION != nodeType((*pNode))) {
164,292✔
1970
    return DEAL_RES_CONTINUE;
135,884✔
1971
  }
1972
  SFunctionNode* pFuncNode = *(SFunctionNode**)(pNode);
28,408✔
1973
  if (!fmIsStreamPesudoColVal(pFuncNode->funcId)) {
28,408✔
1974
    return DEAL_RES_CONTINUE;
884✔
1975
  }
1976
  pData->code = fmSetStreamPseudoFuncParamVal(pFuncNode->funcId, pFuncNode->pParameterList, pData->pStreamRuntimeInfo);
27,164✔
1977
  if (pData->code != TSDB_CODE_SUCCESS) {
27,164✔
1978
    return DEAL_RES_ERROR;
×
1979
  }
1980
  SNode* pFirstParam = nodesListGetNode(pFuncNode->pParameterList, 0);
27,164✔
1981
  ((SValueNode*)pFirstParam)->translate = true;
27,164✔
1982
  SValueNode* res = NULL;
27,164✔
1983
  pData->code = nodesCloneNode(pFirstParam, (SNode**)&res);
27,164✔
1984
  if (NULL == res) {
27,524✔
1985
    return DEAL_RES_ERROR;
×
1986
  }
1987
  nodesDestroyNode(*pNode);
27,524✔
1988
  *pNode = (SNode*)res;
27,164✔
1989

1990
  return DEAL_RES_CONTINUE;
27,164✔
1991
}
1992

1993
static void extractTagColId(SOperatorNode* pOpNode, SArray* pColIdArray) {
45,360✔
1994
  SNode* pLeft = pOpNode->pLeft;
45,360✔
1995
  SNode* pRight = pOpNode->pRight;
45,360✔
1996
  SColumnNode* pColNode = nodeType(pLeft) == QUERY_NODE_COLUMN ?
45,360✔
1997
    (SColumnNode*)pLeft : (SColumnNode*)pRight;
45,360✔
1998

1999
  col_id_t colId = pColNode->colId;
45,360✔
2000
  void* _tmp = taosArrayPush(pColIdArray, &colId);
45,360✔
2001
}
45,360✔
2002

2003
static int32_t buildTagCondKey(
22,680✔
2004
  const SNode* pTagCond, char** pTagCondKey,
2005
  int32_t* tagCondKeyLen, SArray** pTagColIds) {
2006
  if (NULL == pTagCond ||
22,680✔
2007
    (nodeType(pTagCond) != QUERY_NODE_OPERATOR &&
22,680✔
2008
      nodeType(pTagCond) != QUERY_NODE_LOGIC_CONDITION)) {
22,680✔
2009
    qError("invalid parameter to extract tag filter symbol");
×
2010
    return TSDB_CODE_INTERNAL_ERROR;
×
2011
  }
2012
  int32_t code = TSDB_CODE_SUCCESS;
22,680✔
2013
  int32_t lino = 0;
22,680✔
2014
  *pTagColIds = taosArrayInit(4, sizeof(col_id_t));
22,680✔
2015

2016
  if (nodeType(pTagCond) == QUERY_NODE_OPERATOR) {
22,680✔
2017
    extractTagColId((SOperatorNode*)pTagCond, *pTagColIds);
×
2018
  } else if (nodeType(pTagCond) == QUERY_NODE_LOGIC_CONDITION) {
22,680✔
2019
    SNode* pChild = NULL;
22,680✔
2020
    FOREACH(pChild, ((SLogicConditionNode*)pTagCond)->pParameterList) {
68,040✔
2021
      extractTagColId((SOperatorNode*)pChild, *pTagColIds);
45,360✔
2022
    }
2023
  }
2024

2025
  taosArraySort(*pTagColIds, compareUint16Val);
22,680✔
2026

2027
  // encode ordered colIds into key string, separated by ','
2028
  *tagCondKeyLen =
45,360✔
2029
    (int32_t)(taosArrayGetSize(*pTagColIds) * (sizeof(col_id_t) + 1) - 1);
22,680✔
2030
  *pTagCondKey = (char*)taosMemoryCalloc(1, *tagCondKeyLen);
22,680✔
2031
  TSDB_CHECK_NULL(*pTagCondKey, code, lino, _end, terrno);
22,680✔
2032
  char* pStart = *pTagCondKey;
22,680✔
2033
  for (int32_t i = 0; i < taosArrayGetSize(*pTagColIds); ++i) {
68,040✔
2034
    col_id_t* pColId = (col_id_t*)taosArrayGet(*pTagColIds, i);
45,360✔
2035
    TSDB_CHECK_NULL(pColId, code, lino, _end, terrno);
45,360✔
2036
    memcpy(pStart, pColId, sizeof(col_id_t));
45,360✔
2037
    pStart += sizeof(col_id_t);
45,360✔
2038
    if (i != taosArrayGetSize(*pTagColIds) - 1) {
45,360✔
2039
      *pStart = ',';
22,680✔
2040
      pStart += 1;
22,680✔
2041
    }
2042
  }
2043

2044
_end:
22,680✔
2045
  if (TSDB_CODE_SUCCESS != code) {
22,680✔
2046
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
2047
    terrno = code;
×
2048
  }
2049
  return code;
22,680✔
2050
}
2051

2052
static EDealRes canOptimizeTagCondFilter(SNode* pTagCond, void* pContext) {
210,960✔
2053
  if (NULL == pTagCond) {
210,960✔
2054
    *(bool*)pContext = false;
×
2055
    return DEAL_RES_END;
×
2056
  }
2057
  if (nodeType(pTagCond) == QUERY_NODE_VALUE ||
210,960✔
2058
    nodeType(pTagCond) == QUERY_NODE_COLUMN) {
140,040✔
2059
    return DEAL_RES_CONTINUE;
116,280✔
2060
  }
2061
  if (nodeType(pTagCond) == QUERY_NODE_OPERATOR &&
94,680✔
2062
    ((SOperatorNode*)pTagCond)->opType == OP_TYPE_EQUAL) {
46,440✔
2063
    return DEAL_RES_CONTINUE;
45,360✔
2064
  }
2065
  if (nodeType(pTagCond) == QUERY_NODE_LOGIC_CONDITION &&
49,320✔
2066
    ((SLogicConditionNode*)pTagCond)->condType == LOGIC_COND_TYPE_AND) {
22,680✔
2067
    return DEAL_RES_CONTINUE;
22,680✔
2068
  }
2069
  if (nodeType(pTagCond) == QUERY_NODE_FUNCTION &&
52,200✔
2070
    fmIsStreamPesudoColVal(((SFunctionNode*)pTagCond)->funcId)) {
25,560✔
2071
    return DEAL_RES_CONTINUE;
25,560✔
2072
  }
2073
  *(bool*)pContext = false;
1,080✔
2074
  return DEAL_RES_END;
1,080✔
2075
}
2076

2077
int32_t getTableList(void* pVnode, SScanPhysiNode* pScanNode, SNode* pTagCond, SNode* pTagIndexCond,
219,647,338✔
2078
                     STableListInfo* pListInfo, uint8_t* digest, const char* idstr, SStorageAPI* pStorageAPI, void* pStreamInfo) {
2079
  int32_t code = TSDB_CODE_SUCCESS;
219,647,338✔
2080
  int32_t lino = 0;
219,647,338✔
2081
  size_t  numOfTables = 0;
219,647,338✔
2082
  bool    listAdded = false;
219,647,338✔
2083

2084
  pListInfo->idInfo.suid = pScanNode->suid;
219,735,063✔
2085
  pListInfo->idInfo.tableType = pScanNode->tableType;
219,650,885✔
2086

2087
  SArray* pUidList = taosArrayInit(8, sizeof(uint64_t));
219,558,486✔
2088
  QUERY_CHECK_NULL(pUidList, code, lino, _error, terrno);
219,471,629✔
2089

2090
  SIdxFltStatus status = SFLT_NOT_INDEX;
219,471,629✔
2091
  char*   pTagCondKey = NULL;
219,429,489✔
2092
  int32_t tagCondKeyLen;
219,559,298✔
2093
  SArray* pTagColIds = NULL;
219,475,080✔
2094
  char*   pPayload = NULL;
219,579,074✔
2095
  qTrace("getTableList called, suid:%" PRIu64
219,579,074✔
2096
    ", tagCond:%p, tagIndexCond:%p, %d %d", pScanNode->suid, pTagCond,
2097
    pTagIndexCond, pScanNode->tableType, pScanNode->virtualStableScan);
2098
  if (pScanNode->tableType != TSDB_SUPER_TABLE && !pScanNode->virtualStableScan) {
219,579,074✔
2099
    pListInfo->idInfo.uid = pScanNode->uid;
148,802,575✔
2100
    if (pStorageAPI->metaFn.isTableExisted(pVnode, pScanNode->uid)) {
148,792,005✔
2101
      void* tmp = taosArrayPush(pUidList, &pScanNode->uid);
148,900,519✔
2102
      QUERY_CHECK_NULL(tmp, code, lino, _error, terrno);
148,884,269✔
2103
    }
2104
    code = doFilterByTagCond(pListInfo, pUidList, pTagCond, pVnode, status, pStorageAPI, false, &listAdded, pStreamInfo);
148,938,594✔
2105
    QUERY_CHECK_CODE(code, lino, _end);
148,918,569✔
2106
  } else {
2107
    bool      isStream = (pStreamInfo != NULL);
70,933,438✔
2108
    bool      hasTagCond = (pTagCond != NULL);
70,933,438✔
2109
    bool      canCacheTagEqCondFilter = false;
70,933,438✔
2110
    T_MD5_CTX context = {0};
70,835,127✔
2111

2112
    qTrace("start to get table list by tag filter, suid:%" PRIu64
70,909,667✔
2113
      ",tsStableTagFilterCache:%d, tsTagFilterCache:%d", 
2114
      pScanNode->suid, tsStableTagFilterCache, tsTagFilterCache);
2115

2116
    bool acquired = false;
70,909,667✔
2117
    // first, check whether we can use stable tag filter cache
2118
    if (tsStableTagFilterCache && isStream && hasTagCond) {
70,824,968✔
2119
      canCacheTagEqCondFilter = true;
23,760✔
2120
      nodesWalkExpr(pTagCond, canOptimizeTagCondFilter,
23,760✔
2121
        (void*)&canCacheTagEqCondFilter);
2122
    }
2123
    if (canCacheTagEqCondFilter) {
70,745,782✔
2124
      qDebug("%s, stable tag filter condition can be optimized", idstr);
22,680✔
2125
      if (((SStreamRuntimeFuncInfo*)pStreamInfo)->hasPlaceHolder) {
22,680✔
2126
        SNode* tmp = NULL;
22,680✔
2127
        code = nodesCloneNode((SNode*)pTagCond, &tmp);
22,680✔
2128
        QUERY_CHECK_CODE(code, lino, _error);
22,680✔
2129

2130
        PlaceHolderContext ctx = {.code = TSDB_CODE_SUCCESS, .pStreamRuntimeInfo = (SStreamRuntimeFuncInfo*)pStreamInfo};
22,680✔
2131
        nodesRewriteExpr(&tmp, replacePlaceHolderColumn, (void*)&ctx);
22,680✔
2132
        if (TSDB_CODE_SUCCESS != ctx.code) {
22,680✔
2133
          nodesDestroyNode(tmp);
×
2134
          code = ctx.code;
×
2135
          goto _error;
×
2136
        }
2137
        code = genStableTagFilterDigest(tmp, &context);
22,680✔
2138
        nodesDestroyNode(tmp);
22,680✔
2139
      } else {
2140
        code = genStableTagFilterDigest(pTagCond, &context);
×
2141
      }
2142
      QUERY_CHECK_CODE(code, lino, _error);
22,680✔
2143

2144
      code = buildTagCondKey(
22,680✔
2145
        pTagCond, &pTagCondKey, &tagCondKeyLen, &pTagColIds);
2146
      QUERY_CHECK_CODE(code, lino, _error);
22,680✔
2147
      code = pStorageAPI->metaFn.getStableCachedTableList(
22,680✔
2148
        pVnode, pScanNode->suid, pTagCondKey, tagCondKeyLen,
22,680✔
2149
        context.digest, tListLen(context.digest), pUidList, &acquired);
2150
      QUERY_CHECK_CODE(code, lino, _error);
22,680✔
2151
    } else if (tsTagFilterCache) {
70,723,102✔
2152
      // second, try to use normal tag filter cache
2153
      qDebug("%s using normal tag filter cache", idstr);
63,396✔
2154
      if (pStreamInfo != NULL && ((SStreamRuntimeFuncInfo*)pStreamInfo)->hasPlaceHolder) {
65,360✔
2155
        SNode* tmp = NULL;
1,964✔
2156
        code = nodesCloneNode((SNode*)pTagCond, &tmp);
1,964✔
2157
        QUERY_CHECK_CODE(code, lino, _error);
1,964✔
2158

2159
        PlaceHolderContext ctx = {.code = TSDB_CODE_SUCCESS, .pStreamRuntimeInfo = (SStreamRuntimeFuncInfo*)pStreamInfo};
1,964✔
2160
        nodesRewriteExpr(&tmp, replacePlaceHolderColumn, (void*)&ctx);
1,964✔
2161
        if (TSDB_CODE_SUCCESS != ctx.code) {
1,964✔
2162
          nodesDestroyNode(tmp);
×
2163
          code = ctx.code;
×
2164
          goto _error;
×
2165
        }
2166
        code = genTagFilterDigest(tmp, &context);
1,964✔
2167
        nodesDestroyNode(tmp);
1,964✔
2168
      } else {
2169
        code = genTagFilterDigest(pTagCond, &context);
61,432✔
2170
      }
2171
      // try to retrieve the result from meta cache
2172
      QUERY_CHECK_CODE(code, lino, _error);      
63,396✔
2173
      code = pStorageAPI->metaFn.getCachedTableList(
63,396✔
2174
        pVnode, pScanNode->suid, context.digest,
63,396✔
2175
        tListLen(context.digest), pUidList, &acquired);
2176
      QUERY_CHECK_CODE(code, lino, _error);
6,857✔
2177
    }
2178
    if (acquired) {
70,651,111✔
2179
      taosArrayDestroy(pTagColIds);
60,580✔
2180
      pTagColIds = NULL;
60,580✔
2181
      
2182
      digest[0] = 1;
60,580✔
2183
      memcpy(
121,160✔
2184
        digest + 1, context.digest, tListLen(context.digest));
60,580✔
2185
      qDebug("suid:%" PRIu64 ", %s retrieve table uid list from cache,"
60,580✔
2186
        " numOfTables:%d", 
2187
        pScanNode->suid, idstr, (int32_t)taosArrayGetSize(pUidList));
2188
      goto _end;
60,580✔
2189
    } else {
2190
      qDebug("suid:%" PRIu64 
70,590,531✔
2191
        ", failed to get table uid list from cache", pScanNode->suid);
2192
    }
2193

2194
    if (!pTagCond) {  // no tag filter condition exists, let's fetch all tables of this super table
70,845,733✔
2195
      code = pStorageAPI->metaFn.getChildTableList(pVnode, pScanNode->suid, pUidList);
57,983,567✔
2196
      QUERY_CHECK_CODE(code, lino, _error);
57,957,129✔
2197
      qTrace("no tag filter, get all child tables, numOfTables:%d", (int32_t)taosArrayGetSize(pUidList));
57,957,129✔
2198
    } else {
2199
      // failed to find the result in the cache, let try to calculate the results
2200
      if (pTagIndexCond) {
12,862,166✔
2201
        void* pIndex = pStorageAPI->metaFn.getInvertIndex(pVnode);
4,481,728✔
2202

2203
        SIndexMetaArg metaArg = {.metaEx = pVnode,
4,481,796✔
2204
                                 .idx = pStorageAPI->metaFn.storeGetIndexInfo(pVnode),
4,481,728✔
2205
                                 .ivtIdx = pIndex,
2206
                                 .suid = pScanNode->uid};
4,481,259✔
2207

2208
        status = SFLT_NOT_INDEX;
4,481,259✔
2209
        code = doFilterTag(pTagIndexCond, &metaArg, pUidList, &status, &pStorageAPI->metaFilter);
4,481,259✔
2210
        if (code != 0 || status == SFLT_NOT_INDEX) {  // temporarily disable it for performance sake
4,476,047✔
2211
          qDebug("failed to get tableIds from index, suid:%" PRIu64 ", uidListSize:%d", pScanNode->uid, (int32_t)taosArrayGetSize(pUidList));
1,066,665✔
2212
        } else {
2213
          qDebug("succ to get filter result, table num: %d", (int)taosArrayGetSize(pUidList));
3,409,382✔
2214
        }
2215
      }
2216
    }
2217
    qTrace("after index filter, pTagCond:%p uidListSize:%d", pTagCond, (int32_t)taosArrayGetSize(pUidList));
70,814,481✔
2218
    code = doFilterByTagCond(pListInfo, pUidList, pTagCond, pVnode, status,
70,816,898✔
2219
      pStorageAPI, tsTagFilterCache || tsStableTagFilterCache,
70,816,898✔
2220
      &listAdded, pStreamInfo);
2221
    QUERY_CHECK_CODE(code, lino, _error);
70,801,628✔
2222

2223
    // let's add the filter results into meta-cache
2224
    numOfTables = taosArrayGetSize(pUidList);
70,800,526✔
2225

2226
    if (canCacheTagEqCondFilter) {
70,805,469✔
2227
      qInfo("%s, suid:%" PRIu64 ", add uid list to stable tag filter cache, "
10,080✔
2228
            "uidListSize:%d, origin key:%" PRIu64 ",%" PRIu64,
2229
            idstr, pScanNode->suid, (int32_t)numOfTables,
2230
            *(uint64_t*)context.digest, *(uint64_t*)(context.digest + 8));
2231

2232
      code = pStorageAPI->metaFn.putStableCachedTableList(
10,080✔
2233
        pVnode, pScanNode->suid, pTagCondKey, tagCondKeyLen,
2234
        context.digest, tListLen(context.digest),
2235
        pUidList, &pTagColIds);
2236
      QUERY_CHECK_CODE(code, lino, _end);
10,080✔
2237

2238
      digest[0] = 1;
10,080✔
2239
      memcpy(digest + 1, context.digest, tListLen(context.digest));
10,080✔
2240
    } else if (tsTagFilterCache) {
70,795,389✔
2241
      qInfo("%s, suid:%" PRIu64 ", add uid list to normal tag filter cache, "
15,416✔
2242
            "uidListSize:%d, origin key:%" PRIu64 ",%" PRIu64,
2243
            idstr, pScanNode->suid, (int32_t)numOfTables,
2244
            *(uint64_t*)context.digest, *(uint64_t*)(context.digest + 8));
2245
      size_t size = numOfTables * sizeof(uint64_t) + sizeof(int32_t);
15,416✔
2246
      pPayload = taosMemoryMalloc(size);
15,416✔
2247
      QUERY_CHECK_NULL(pPayload, code, lino, _end, terrno);
15,416✔
2248

2249
      *(int32_t*)pPayload = (int32_t)numOfTables;
15,416✔
2250
      if (numOfTables > 0) {
15,416✔
2251
        void* tmp = taosArrayGet(pUidList, 0);
12,732✔
2252
        QUERY_CHECK_NULL(tmp, code, lino, _end, terrno);
12,732✔
2253
        memcpy(pPayload + sizeof(int32_t), tmp, numOfTables * sizeof(uint64_t));
12,732✔
2254
      }
2255

2256
      code = pStorageAPI->metaFn.putCachedTableList(pVnode, pScanNode->suid,
15,416✔
2257
                                                    context.digest,
2258
                                                    tListLen(context.digest),
2259
                                                    pPayload, size, 1);
2260
      if (TSDB_CODE_SUCCESS == code) {
15,416✔
2261
        /*
2262
          data referenced by pPayload is used in lru cache,
2263
          reset pPayload to NULL to avoid being freed in _error block
2264
        */
2265
        pPayload = NULL;
15,056✔
2266
      } else {
2267
        if (TSDB_CODE_DUP_KEY == code) {
360✔
2268
          /*
2269
            another thread has already put the same key into cache,
2270
            we can just ignore this error
2271
          */
2272
          code = TSDB_CODE_SUCCESS;
360✔
2273
        }
2274
        QUERY_CHECK_CODE(code, lino, _end);
360✔
2275
      }
2276

2277

2278
      digest[0] = 1;
15,416✔
2279
      memcpy(digest + 1, context.digest, tListLen(context.digest));
15,416✔
2280
    }
2281
  }
2282

2283
_end:
219,773,632✔
2284
  if (!listAdded) {
219,824,567✔
2285
    numOfTables = taosArrayGetSize(pUidList);
208,614,727✔
2286
    for (int i = 0; i < numOfTables; i++) {
630,062,950✔
2287
      void* tmp = taosArrayGet(pUidList, i);
421,428,457✔
2288
      QUERY_CHECK_NULL(tmp, code, lino, _error, terrno);
421,474,301✔
2289
      STableKeyInfo info = {.uid = *(uint64_t*)tmp, .groupId = 0};
421,474,301✔
2290

2291
      void* p = taosArrayPush(pListInfo->pTableList, &info);
421,455,552✔
2292
      if (p == NULL) {
421,549,537✔
2293
        taosArrayDestroy(pUidList);
×
2294
        return terrno;
×
2295
      }
2296

2297
      qTrace("tagfilter get uid:%" PRIu64 ", %s", info.uid, idstr);
421,549,537✔
2298
    }
2299
  }
2300

2301
  qDebug("%s, table list with %d uids built", idstr, (int32_t)numOfTables);
219,844,333✔
2302

2303
_error:
219,840,682✔
2304
  taosArrayDestroy(pUidList);
219,840,724✔
2305
  taosArrayDestroy(pTagColIds);
219,745,571✔
2306
  taosMemFreeClear(pTagCondKey);
219,763,016✔
2307
  taosMemFreeClear(pPayload);
219,763,016✔
2308
  if (code != TSDB_CODE_SUCCESS) {
219,763,016✔
2309
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
1,102✔
2310
  }
2311
  return code;
219,773,262✔
2312
}
2313

2314
int32_t qGetTableList(int64_t suid, void* pVnode, void* node, SArray** tableList, void* pTaskInfo) {
4,294✔
2315
  int32_t        code = TSDB_CODE_SUCCESS;
4,294✔
2316
  int32_t        lino = 0;
4,294✔
2317
  SSubplan*      pSubplan = (SSubplan*)node;
4,294✔
2318
  SScanPhysiNode pNode = {0};
4,294✔
2319
  pNode.suid = suid;
4,294✔
2320
  pNode.uid = suid;
4,294✔
2321
  pNode.tableType = TSDB_SUPER_TABLE;
4,294✔
2322

2323
  STableListInfo* pTableListInfo = tableListCreate();
4,294✔
2324
  QUERY_CHECK_NULL(pTableListInfo, code, lino, _end, terrno);
4,294✔
2325
  uint8_t digest[17] = {0};
4,294✔
2326
  code = getTableList(pVnode, &pNode, pSubplan ? pSubplan->pTagCond : NULL, pSubplan ? pSubplan->pTagIndexCond : NULL,
4,294✔
2327
                      pTableListInfo, digest, "qGetTableList", &((SExecTaskInfo*)pTaskInfo)->storageAPI, NULL);
2328
  QUERY_CHECK_CODE(code, lino, _end);
4,294✔
2329
  *tableList = pTableListInfo->pTableList;
4,294✔
2330
  pTableListInfo->pTableList = NULL;
4,294✔
2331
  tableListDestroy(pTableListInfo);
4,294✔
2332

2333
_end:
4,294✔
2334
  if (code != TSDB_CODE_SUCCESS) {
4,294✔
2335
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
2336
  }
2337
  return code;
4,294✔
2338
}
2339

2340
size_t getTableTagsBufLen(const SNodeList* pGroups) {
×
2341
  size_t keyLen = 0;
×
2342

2343
  SNode* node;
2344
  FOREACH(node, pGroups) {
×
2345
    SExprNode* pExpr = (SExprNode*)node;
×
2346
    keyLen += pExpr->resType.bytes;
×
2347
  }
2348

2349
  keyLen += sizeof(int8_t) * LIST_LENGTH(pGroups);
×
2350
  return keyLen;
×
2351
}
2352

2353
int32_t getGroupIdFromTagsVal(void* pVnode, uint64_t uid, SNodeList* pGroupNode, char* keyBuf, uint64_t* pGroupId,
×
2354
                              SStorageAPI* pAPI) {
2355
  SMetaReader mr = {0};
×
2356

2357
  pAPI->metaReaderFn.initReader(&mr, pVnode, META_READER_LOCK, &pAPI->metaFn);
×
2358
  if (pAPI->metaReaderFn.getEntryGetUidCache(&mr, uid) != 0) {  // table not exist
×
2359
    pAPI->metaReaderFn.clearReader(&mr);
×
2360
    return TSDB_CODE_PAR_TABLE_NOT_EXIST;
×
2361
  }
2362

2363
  SNodeList* groupNew = NULL;
×
2364
  int32_t    code = nodesCloneList(pGroupNode, &groupNew);
×
2365
  if (TSDB_CODE_SUCCESS != code) {
×
2366
    pAPI->metaReaderFn.clearReader(&mr);
×
2367
    return code;
×
2368
  }
2369

2370
  STransTagExprCtx ctx = {.code = 0, .pReader = &mr};
×
2371
  nodesRewriteExprsPostOrder(groupNew, doTranslateTagExpr, &ctx);
×
2372
  if (TSDB_CODE_SUCCESS != ctx.code) {
×
2373
    nodesDestroyList(groupNew);
×
2374
    pAPI->metaReaderFn.clearReader(&mr);
×
2375
    return code;
×
2376
  }
2377
  char* isNull = (char*)keyBuf;
×
2378
  char* pStart = (char*)keyBuf + sizeof(int8_t) * LIST_LENGTH(pGroupNode);
×
2379

2380
  SNode*  pNode;
2381
  int32_t index = 0;
×
2382
  FOREACH(pNode, groupNew) {
×
2383
    SNode*  pNew = NULL;
×
2384
    int32_t code = scalarCalculateConstants(pNode, &pNew);
×
2385
    if (TSDB_CODE_SUCCESS == code) {
×
2386
      REPLACE_NODE(pNew);
×
2387
    } else {
2388
      nodesDestroyList(groupNew);
×
2389
      pAPI->metaReaderFn.clearReader(&mr);
×
2390
      return code;
×
2391
    }
2392

2393
    if (nodeType(pNew) != QUERY_NODE_VALUE) {
×
2394
      nodesDestroyList(groupNew);
×
2395
      pAPI->metaReaderFn.clearReader(&mr);
×
2396
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(TSDB_CODE_QRY_EXECUTOR_INTERNAL_ERROR));
×
2397
      return TSDB_CODE_QRY_EXECUTOR_INTERNAL_ERROR;
×
2398
    }
2399
    SValueNode* pValue = (SValueNode*)pNew;
×
2400

2401
    if (pValue->node.resType.type == TSDB_DATA_TYPE_NULL || pValue->isNull) {
×
2402
      isNull[index++] = 1;
×
2403
      continue;
×
2404
    } else {
2405
      isNull[index++] = 0;
×
2406
      char* data = nodesGetValueFromNode(pValue);
×
2407
      if (pValue->node.resType.type == TSDB_DATA_TYPE_JSON) {
×
2408
        if (tTagIsJson(data)) {
×
2409
          terrno = TSDB_CODE_QRY_JSON_IN_GROUP_ERROR;
×
2410
          nodesDestroyList(groupNew);
×
2411
          pAPI->metaReaderFn.clearReader(&mr);
×
2412
          return terrno;
×
2413
        }
2414
        int32_t len = getJsonValueLen(data);
×
2415
        memcpy(pStart, data, len);
×
2416
        pStart += len;
×
2417
      } else if (IS_VAR_DATA_TYPE(pValue->node.resType.type)) {
×
2418
        if (IS_STR_DATA_BLOB(pValue->node.resType.type)) {
×
2419
          return TSDB_CODE_BLOB_NOT_SUPPORT_TAG;
×
2420
        }
2421
        memcpy(pStart, data, varDataTLen(data));
×
2422
        pStart += varDataTLen(data);
×
2423
      } else {
2424
        memcpy(pStart, data, pValue->node.resType.bytes);
×
2425
        pStart += pValue->node.resType.bytes;
×
2426
      }
2427
    }
2428
  }
2429

2430
  int32_t len = (int32_t)(pStart - (char*)keyBuf);
×
2431
  *pGroupId = calcGroupId(keyBuf, len);
×
2432

2433
  nodesDestroyList(groupNew);
×
2434
  pAPI->metaReaderFn.clearReader(&mr);
×
2435

2436
  return TSDB_CODE_SUCCESS;
×
2437
}
2438

2439
SArray* makeColumnArrayFromList(SNodeList* pNodeList) {
8,410,727✔
2440
  if (!pNodeList) {
8,410,727✔
2441
    return NULL;
×
2442
  }
2443

2444
  size_t  numOfCols = LIST_LENGTH(pNodeList);
8,410,727✔
2445
  SArray* pList = taosArrayInit(numOfCols, sizeof(SColumn));
8,410,583✔
2446
  if (pList == NULL) {
8,409,947✔
2447
    return NULL;
×
2448
  }
2449

2450
  for (int32_t i = 0; i < numOfCols; ++i) {
18,815,849✔
2451
    SColumnNode* pColNode = (SColumnNode*)nodesListGetNode(pNodeList, i);
10,408,237✔
2452
    if (!pColNode) {
10,409,819✔
2453
      taosArrayDestroy(pList);
×
2454
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(TSDB_CODE_QRY_EXECUTOR_INTERNAL_ERROR));
×
2455
      return NULL;
×
2456
    }
2457

2458
    // todo extract method
2459
    SColumn c = {0};
10,409,819✔
2460
    c.slotId = pColNode->slotId;
10,408,808✔
2461
    c.colId = pColNode->colId;
10,408,356✔
2462
    c.type = pColNode->node.resType.type;
10,408,780✔
2463
    c.bytes = pColNode->node.resType.bytes;
10,408,905✔
2464
    c.precision = pColNode->node.resType.precision;
10,407,919✔
2465
    c.scale = pColNode->node.resType.scale;
10,408,808✔
2466

2467
    void* tmp = taosArrayPush(pList, &c);
10,406,489✔
2468
    if (!tmp) {
10,406,489✔
2469
      taosArrayDestroy(pList);
×
2470
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
2471
      return NULL;
×
2472
    }
2473
  }
2474

2475
  return pList;
8,407,612✔
2476
}
2477

2478
int32_t extractColMatchInfo(SNodeList* pNodeList, SDataBlockDescNode* pOutputNodeList, int32_t* numOfOutputCols,
270,074,336✔
2479
                            int32_t type, SColMatchInfo* pMatchInfo) {
2480
  size_t  numOfCols = LIST_LENGTH(pNodeList);
270,074,336✔
2481
  int32_t code = TSDB_CODE_SUCCESS;
270,083,608✔
2482
  int32_t lino = 0;
270,083,608✔
2483

2484
  pMatchInfo->matchType = type;
270,083,608✔
2485

2486
  SArray* pList = taosArrayInit(numOfCols, sizeof(SColMatchItem));
270,072,694✔
2487
  if (pList == NULL) {
269,910,996✔
2488
    code = terrno;
×
2489
    return code;
×
2490
  }
2491

2492
  for (int32_t i = 0; i < numOfCols; ++i) {
1,244,796,312✔
2493
    STargetNode* pNode = (STargetNode*)nodesListGetNode(pNodeList, i);
974,759,200✔
2494
    QUERY_CHECK_NULL(pNode, code, lino, _end, terrno);
974,912,250✔
2495
    if (nodeType(pNode->pExpr) == QUERY_NODE_COLUMN) {
974,912,250✔
2496
      SColumnNode* pColNode = (SColumnNode*)pNode->pExpr;
969,378,684✔
2497

2498
      SColMatchItem c = {.needOutput = true};
969,436,617✔
2499
      c.colId = pColNode->colId;
969,415,014✔
2500
      c.srcSlotId = pColNode->slotId;
969,334,763✔
2501
      c.dstSlotId = pNode->slotId;
969,311,357✔
2502
      c.isPk = pColNode->isPk;
969,449,602✔
2503
      c.dataType = pColNode->node.resType;
969,363,181✔
2504
      void* tmp = taosArrayPush(pList, &c);
969,390,008✔
2505
      QUERY_CHECK_NULL(tmp, code, lino, _end, terrno);
969,390,008✔
2506
    }
2507
  }
2508

2509
  // set the output flag for each column in SColMatchInfo, according to the
2510
  *numOfOutputCols = 0;
270,037,112✔
2511
  int32_t num = LIST_LENGTH(pOutputNodeList->pSlots);
270,166,363✔
2512
  for (int32_t i = 0; i < num; ++i) {
1,349,199,096✔
2513
    SSlotDescNode* pNode = (SSlotDescNode*)nodesListGetNode(pOutputNodeList->pSlots, i);
1,079,116,143✔
2514
    QUERY_CHECK_NULL(pNode, code, lino, _end, terrno);
1,079,140,392✔
2515

2516
    // todo: add reserve flag check
2517
    // it is a column reserved for the arithmetic expression calculation
2518
    if (pNode->slotId >= numOfCols) {
1,079,140,392✔
2519
      (*numOfOutputCols) += 1;
104,345,961✔
2520
      continue;
104,346,675✔
2521
    }
2522

2523
    SColMatchItem* info = NULL;
974,889,083✔
2524
    for (int32_t j = 0; j < taosArrayGetSize(pList); ++j) {
2,147,483,647✔
2525
      info = taosArrayGet(pList, j);
2,147,483,647✔
2526
      QUERY_CHECK_NULL(info, code, lino, _end, terrno);
2,147,483,647✔
2527
      if (info->dstSlotId == pNode->slotId) {
2,147,483,647✔
2528
        break;
968,497,075✔
2529
      }
2530
    }
2531

2532
    if (pNode->output) {
12,179,721✔
2533
      (*numOfOutputCols) += 1;
965,794,450✔
2534
    } else if (info != NULL) {
9,035,797✔
2535
      // select distinct tbname from stb where tbname='abc';
2536
      info->needOutput = false;
9,049,646✔
2537
    }
2538
  }
2539

2540
  pMatchInfo->pList = pList;
270,082,953✔
2541

2542
_end:
270,127,511✔
2543
  if (code != TSDB_CODE_SUCCESS) {
270,127,511✔
2544
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
2545
  }
2546
  return code;
270,082,868✔
2547
}
2548

2549
static SResSchema createResSchema(int32_t type, int32_t bytes, int32_t slotId, int32_t scale, int32_t precision,
815,911,155✔
2550
                                  const char* name) {
2551
  SResSchema s = {0};
815,911,155✔
2552
  s.scale = scale;
815,993,712✔
2553
  s.type = type;
815,993,712✔
2554
  s.bytes = bytes;
815,993,712✔
2555
  s.slotId = slotId;
815,993,712✔
2556
  s.precision = precision;
815,993,712✔
2557
  tstrncpy(s.name, name, tListLen(s.name));
815,993,712✔
2558

2559
  return s;
815,993,712✔
2560
}
2561

2562
static SColumn* createColumn(int32_t blockId, int32_t slotId, int32_t colId, SDataType* pType, EColumnType colType) {
737,489,649✔
2563
  SColumn* pCol = taosMemoryCalloc(1, sizeof(SColumn));
737,489,649✔
2564
  if (pCol == NULL) {
736,956,442✔
2565
    return NULL;
×
2566
  }
2567

2568
  pCol->slotId = slotId;
736,956,442✔
2569
  pCol->colId = colId;
737,006,921✔
2570
  pCol->bytes = pType->bytes;
737,103,138✔
2571
  pCol->type = pType->type;
737,198,008✔
2572
  pCol->scale = pType->scale;
737,425,063✔
2573
  pCol->precision = pType->precision;
737,353,686✔
2574
  pCol->dataBlockId = blockId;
737,556,970✔
2575
  pCol->colType = colType;
737,519,026✔
2576
  return pCol;
737,514,664✔
2577
}
2578

2579
int32_t createExprFromOneNode(SExprInfo* pExp, SNode* pNode, int16_t slotId) {
820,647,374✔
2580
  int32_t code = TSDB_CODE_SUCCESS;
820,647,374✔
2581
  int32_t lino = 0;
820,647,374✔
2582
  pExp->base.numOfParams = 0;
820,647,374✔
2583
  pExp->base.pParam = NULL;
820,748,377✔
2584
  pExp->pExpr = taosMemoryCalloc(1, sizeof(tExprNode));
820,640,431✔
2585
  QUERY_CHECK_NULL(pExp->pExpr, code, lino, _end, terrno);
819,997,481✔
2586

2587
  pExp->pExpr->_function.num = 1;
820,207,504✔
2588
  pExp->pExpr->_function.functionId = -1;
820,377,302✔
2589

2590
  int32_t type = nodeType(pNode);
820,392,345✔
2591
  // it is a project query, or group by column
2592
  if (type == QUERY_NODE_COLUMN) {
820,550,245✔
2593
    pExp->pExpr->nodeType = QUERY_NODE_COLUMN;
489,464,032✔
2594
    SColumnNode* pColNode = (SColumnNode*)pNode;
489,505,580✔
2595

2596
    pExp->base.pParam = taosMemoryCalloc(1, sizeof(SFunctParam));
489,505,580✔
2597
    QUERY_CHECK_NULL(pExp->base.pParam, code, lino, _end, terrno);
489,251,753✔
2598

2599
    pExp->base.numOfParams = 1;
489,324,034✔
2600

2601
    SDataType* pType = &pColNode->node.resType;
489,315,958✔
2602
    pExp->base.resSchema =
2603
        createResSchema(pType->type, pType->bytes, slotId, pType->scale, pType->precision, pColNode->colName);
489,434,124✔
2604

2605
    pExp->base.pParam[0].pCol =
978,828,545✔
2606
        createColumn(pColNode->dataBlockId, pColNode->slotId, pColNode->colId, pType, pColNode->colType);
978,820,381✔
2607
    QUERY_CHECK_NULL(pExp->base.pParam[0].pCol, code, lino, _end, terrno);
489,479,057✔
2608

2609
    pExp->base.pParam[0].type = FUNC_PARAM_TYPE_COLUMN;
489,264,236✔
2610
  } else if (type == QUERY_NODE_VALUE) {
331,086,213✔
2611
    pExp->pExpr->nodeType = QUERY_NODE_VALUE;
16,389,779✔
2612
    SValueNode* pValNode = (SValueNode*)pNode;
16,389,907✔
2613

2614
    pExp->base.pParam = taosMemoryCalloc(1, sizeof(SFunctParam));
16,389,907✔
2615
    QUERY_CHECK_NULL(pExp->base.pParam, code, lino, _end, terrno);
16,382,687✔
2616

2617
    pExp->base.numOfParams = 1;
16,387,063✔
2618

2619
    SDataType* pType = &pValNode->node.resType;
16,385,700✔
2620
    pExp->base.resSchema =
2621
        createResSchema(pType->type, pType->bytes, slotId, pType->scale, pType->precision, pValNode->node.aliasName);
16,386,331✔
2622
    pExp->base.pParam[0].type = FUNC_PARAM_TYPE_VALUE;
16,385,573✔
2623
    code = nodesValueNodeToVariant(pValNode, &pExp->base.pParam[0].param);
16,381,893✔
2624
    QUERY_CHECK_CODE(code, lino, _end);
16,384,840✔
2625
  } else if (type == QUERY_NODE_REMOTE_VALUE) {
314,696,434✔
2626
    SRemoteValueNode* pRemote = (SRemoteValueNode*)pNode;
27,829,737✔
2627
    code = qFetchRemoteValue(gTaskScalarExtra.pSubJobCtx, pRemote->subQIdx, pRemote);
27,829,737✔
2628
    QUERY_CHECK_CODE(code, lino, _end);
27,849,078✔
2629

2630
    pExp->pExpr->nodeType = QUERY_NODE_VALUE;
23,137,404✔
2631
    SValueNode* pValNode = (SValueNode*)pNode;
23,137,421✔
2632

2633
    pExp->base.pParam = taosMemoryCalloc(1, sizeof(SFunctParam));
23,137,421✔
2634
    QUERY_CHECK_NULL(pExp->base.pParam, code, lino, _end, terrno);
23,137,961✔
2635

2636
    pExp->base.numOfParams = 1;
23,137,401✔
2637

2638
    SDataType* pType = &pValNode->node.resType;
23,137,944✔
2639
    pExp->base.resSchema =
2640
        createResSchema(pType->type, pType->bytes, slotId, pType->scale, pType->precision, pValNode->node.aliasName);
23,137,944✔
2641
    pExp->base.pParam[0].type = FUNC_PARAM_TYPE_VALUE;
23,136,304✔
2642
    code = nodesValueNodeToVariant(pValNode, &pExp->base.pParam[0].param);
23,136,864✔
2643
    QUERY_CHECK_CODE(code, lino, _end);
23,136,844✔
2644
  } else if (type == QUERY_NODE_FUNCTION) {
286,866,697✔
2645
    pExp->pExpr->nodeType = QUERY_NODE_FUNCTION;
255,574,284✔
2646
    SFunctionNode* pFuncNode = (SFunctionNode*)pNode;
255,571,673✔
2647

2648
    SDataType* pType = &pFuncNode->node.resType;
255,571,673✔
2649
    pExp->base.resSchema =
2650
        createResSchema(pType->type, pType->bytes, slotId, pType->scale, pType->precision, pFuncNode->node.aliasName);
255,582,863✔
2651
    tExprNode* pExprNode = pExp->pExpr;
255,568,564✔
2652

2653
    pExprNode->_function.functionId = pFuncNode->funcId;
255,556,859✔
2654
    pExprNode->_function.pFunctNode = pFuncNode;
255,581,014✔
2655
    pExprNode->_function.functionType = pFuncNode->funcType;
255,605,520✔
2656

2657
    tstrncpy(pExprNode->_function.functionName, pFuncNode->functionName, tListLen(pExprNode->_function.functionName));
255,550,763✔
2658

2659
    pExp->base.pParamList = pFuncNode->pParameterList;
255,584,214✔
2660
#if 1
2661
    // todo refactor: add the parameter for tbname function
2662
    const char* name = "tbname";
255,598,163✔
2663
    int32_t     len = strlen(name);
255,598,163✔
2664

2665
    if (!pFuncNode->pParameterList && (memcmp(pExprNode->_function.functionName, name, len) == 0) &&
255,598,163✔
2666
        pExprNode->_function.functionName[len] == 0) {
9,236,002✔
2667
      pFuncNode->pParameterList = NULL;
9,230,408✔
2668
      int32_t     code = nodesMakeList(&pFuncNode->pParameterList);
9,236,899✔
2669
      SValueNode* res = NULL;
9,238,603✔
2670
      if (TSDB_CODE_SUCCESS == code) {
9,239,291✔
2671
        code = nodesMakeNode(QUERY_NODE_VALUE, (SNode**)&res);
9,239,126✔
2672
      }
2673
      QUERY_CHECK_CODE(code, lino, _end);
9,241,571✔
2674
      res->node.resType = (SDataType){.bytes = sizeof(int64_t), .type = TSDB_DATA_TYPE_BIGINT};
9,241,571✔
2675
      code = nodesListAppend(pFuncNode->pParameterList, (SNode*)res);
9,237,463✔
2676
      if (code != TSDB_CODE_SUCCESS) {
9,236,992✔
2677
        nodesDestroyNode((SNode*)res);
×
2678
        res = NULL;
×
2679
      }
2680
      QUERY_CHECK_CODE(code, lino, _end);
9,236,992✔
2681
    }
2682
#endif
2683

2684
    int32_t numOfParam = LIST_LENGTH(pFuncNode->pParameterList);
255,618,347✔
2685

2686
    pExp->base.pParam = taosMemoryCalloc(numOfParam, sizeof(SFunctParam));
255,610,936✔
2687
    QUERY_CHECK_NULL(pExp->base.pParam, code, lino, _end, terrno);
255,514,917✔
2688
    pExp->base.numOfParams = numOfParam;
255,519,569✔
2689

2690
    for (int32_t j = 0; j < numOfParam && TSDB_CODE_SUCCESS == code; ++j) {
631,901,229✔
2691
      SNode* p1 = nodesListGetNode(pFuncNode->pParameterList, j);
376,587,707✔
2692
      QUERY_CHECK_NULL(p1, code, lino, _end, terrno);
376,581,864✔
2693
      if (p1->type == QUERY_NODE_COLUMN) {
376,581,864✔
2694
        SColumnNode* pcn = (SColumnNode*)p1;
247,984,019✔
2695

2696
        pExp->base.pParam[j].type = FUNC_PARAM_TYPE_COLUMN;
247,984,019✔
2697
        pExp->base.pParam[j].pCol =
495,952,148✔
2698
            createColumn(pcn->dataBlockId, pcn->slotId, pcn->colId, &pcn->node.resType, pcn->colType);
495,982,685✔
2699
        QUERY_CHECK_NULL(pExp->base.pParam[j].pCol, code, lino, _end, terrno);
247,987,456✔
2700
      } else if (p1->type == QUERY_NODE_VALUE) {
128,617,441✔
2701
        SValueNode* pvn = (SValueNode*)p1;
64,255,567✔
2702
        pExp->base.pParam[j].type = FUNC_PARAM_TYPE_VALUE;
64,255,567✔
2703
        code = nodesValueNodeToVariant(pvn, &pExp->base.pParam[j].param);
64,245,243✔
2704
        QUERY_CHECK_CODE(code, lino, _end);
64,229,785✔
2705
      } else if (p1->type == QUERY_NODE_REMOTE_VALUE) {
64,396,471✔
2706
        SRemoteValueNode* pRemote = (SRemoteValueNode*)p1;
1,411,146✔
2707
        code = qFetchRemoteValue(gTaskScalarExtra.pSubJobCtx, pRemote->subQIdx, pRemote);
1,411,146✔
2708
        QUERY_CHECK_CODE(code, lino, _end);
1,411,146✔
2709

2710
        SValueNode* pvn = (SValueNode*)pRemote;
1,195,995✔
2711
        pExp->base.pParam[j].type = FUNC_PARAM_TYPE_VALUE;
1,195,995✔
2712
        code = nodesValueNodeToVariant(pvn, &pExp->base.pParam[j].param);
1,195,995✔
2713
        QUERY_CHECK_CODE(code, lino, _end);
1,143,835✔
2714
      }
2715
    }
2716
    pExp->pExpr->_function.bindExprID = ((SExprNode*)pNode)->bindExprID;
255,313,522✔
2717
  } else if (type == QUERY_NODE_OPERATOR) {
31,292,413✔
2718
    pExp->pExpr->nodeType = QUERY_NODE_OPERATOR;
25,605,815✔
2719
    SOperatorNode* pOpNode = (SOperatorNode*)pNode;
25,603,333✔
2720

2721
    pExp->base.pParam = taosMemoryCalloc(1, sizeof(SFunctParam));
25,603,333✔
2722
    QUERY_CHECK_NULL(pExp->base.pParam, code, lino, _end, terrno);
25,600,483✔
2723
    pExp->base.numOfParams = 1;
25,601,394✔
2724

2725
    SDataType* pType = &pOpNode->node.resType;
25,604,074✔
2726
    pExp->base.resSchema =
2727
        createResSchema(pType->type, pType->bytes, slotId, pType->scale, pType->precision, pOpNode->node.aliasName);
25,601,183✔
2728
    pExp->pExpr->_optrRoot.pRootNode = pNode;
25,602,932✔
2729
  } else if (type == QUERY_NODE_CASE_WHEN) {
5,686,994✔
2730
    pExp->pExpr->nodeType = QUERY_NODE_OPERATOR;
5,693,412✔
2731
    SCaseWhenNode* pCaseNode = (SCaseWhenNode*)pNode;
5,693,345✔
2732

2733
    pExp->base.pParam = taosMemoryCalloc(1, sizeof(SFunctParam));
5,693,345✔
2734
    QUERY_CHECK_NULL(pExp->base.pParam, code, lino, _end, terrno);
5,693,412✔
2735
    pExp->base.numOfParams = 1;
5,692,886✔
2736

2737
    SDataType* pType = &pCaseNode->node.resType;
5,692,886✔
2738
    pExp->base.resSchema =
2739
        createResSchema(pType->type, pType->bytes, slotId, pType->scale, pType->precision, pCaseNode->node.aliasName);
5,692,819✔
2740
    pExp->pExpr->_optrRoot.pRootNode = pNode;
5,693,412✔
2741
  } else if (type == QUERY_NODE_LOGIC_CONDITION) {
×
2742
    pExp->pExpr->nodeType = QUERY_NODE_OPERATOR;
1,182✔
2743
    SLogicConditionNode* pCond = (SLogicConditionNode*)pNode;
1,182✔
2744
    pExp->base.pParam = taosMemoryCalloc(1, sizeof(SFunctParam));
1,182✔
2745
    QUERY_CHECK_NULL(pExp->base.pParam, code, lino, _end, terrno);
1,182✔
2746
    pExp->base.numOfParams = 1;
1,182✔
2747
    SDataType* pType = &pCond->node.resType;
1,182✔
2748
    pExp->base.resSchema =
2749
        createResSchema(pType->type, pType->bytes, slotId, pType->scale, pType->precision, pCond->node.aliasName);
1,182✔
2750
    pExp->pExpr->_optrRoot.pRootNode = pNode;
1,182✔
2751
  } else {
2752
    code = TSDB_CODE_QRY_EXECUTOR_INTERNAL_ERROR;
×
2753
    QUERY_CHECK_CODE(code, lino, _end);
×
2754
  }
2755
  pExp->pExpr->relatedTo = ((SExprNode*)pNode)->relatedTo;
815,646,809✔
2756
_end:
820,583,073✔
2757
  if (code != TSDB_CODE_SUCCESS) {
820,583,073✔
2758
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
4,926,825✔
2759
  }
2760
  return code;
820,619,630✔
2761
}
2762

2763
int32_t createExprFromTargetNode(SExprInfo* pExp, STargetNode* pTargetNode) {
820,579,917✔
2764
  return createExprFromOneNode(pExp, pTargetNode->pExpr, pTargetNode->slotId);
820,579,917✔
2765
}
2766

2767
SExprInfo* createExpr(SNodeList* pNodeList, int32_t* numOfExprs) {
×
2768
  *numOfExprs = LIST_LENGTH(pNodeList);
×
2769
  SExprInfo* pExprs = taosMemoryCalloc(*numOfExprs, sizeof(SExprInfo));
×
2770
  if (!pExprs) {
×
2771
    return NULL;
×
2772
  }
2773

2774
  for (int32_t i = 0; i < (*numOfExprs); ++i) {
×
2775
    SExprInfo* pExp = &pExprs[i];
×
2776
    int32_t    code = createExprFromOneNode(pExp, nodesListGetNode(pNodeList, i), i + UD_TAG_COLUMN_INDEX);
×
2777
    if (code != TSDB_CODE_SUCCESS) {
×
2778
      taosMemoryFreeClear(pExprs);
×
2779
      terrno = code;
×
2780
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(code));
×
2781
      return NULL;
×
2782
    }
2783
  }
2784

2785
  return pExprs;
×
2786
}
2787

2788
int32_t createExprInfo(SNodeList* pNodeList, SNodeList* pGroupKeys, SExprInfo** pExprInfo, int32_t* numOfExprs) {
358,638,596✔
2789
  QRY_PARAM_CHECK(pExprInfo);
358,638,596✔
2790

2791
  int32_t code = 0;
358,735,589✔
2792
  int32_t numOfFuncs = LIST_LENGTH(pNodeList);
358,735,589✔
2793
  int32_t numOfGroupKeys = 0;
358,694,903✔
2794
  if (pGroupKeys != NULL) {
358,694,903✔
2795
    numOfGroupKeys = LIST_LENGTH(pGroupKeys);
34,803,391✔
2796
  }
2797

2798
  *numOfExprs = numOfFuncs + numOfGroupKeys;
358,697,190✔
2799
  if (*numOfExprs == 0) {
358,693,628✔
2800
    return code;
43,247,029✔
2801
  }
2802

2803
  SExprInfo* pExprs = taosMemoryCalloc(*numOfExprs, sizeof(SExprInfo));
315,508,711✔
2804
  if (pExprs == NULL) {
315,026,819✔
2805
    return terrno;
×
2806
  }
2807

2808
  for (int32_t i = 0; i < (*numOfExprs); ++i) {
1,130,419,796✔
2809
    STargetNode* pTargetNode = NULL;
820,183,293✔
2810
    if (i < numOfFuncs) {
820,183,293✔
2811
      pTargetNode = (STargetNode*)nodesListGetNode(pNodeList, i);
779,468,309✔
2812
    } else {
2813
      pTargetNode = (STargetNode*)nodesListGetNode(pGroupKeys, i - numOfFuncs);
40,714,984✔
2814
    }
2815
    if (!pTargetNode) {
820,440,578✔
2816
      destroyExprInfo(pExprs, *numOfExprs);
×
2817
      taosMemoryFreeClear(pExprs);
×
2818
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
2819
      return terrno;
×
2820
    }
2821

2822
    SExprInfo* pExp = &pExprs[i];
820,440,578✔
2823
    code = createExprFromTargetNode(pExp, pTargetNode);
820,461,874✔
2824
    if (code != TSDB_CODE_SUCCESS) {
820,319,802✔
2825
      destroyExprInfo(pExprs, *numOfExprs);
4,926,825✔
2826
      taosMemoryFreeClear(pExprs);
4,926,825✔
2827
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(code));
4,926,825✔
2828
      return code;
4,926,825✔
2829
    }
2830
  }
2831

2832
  *pExprInfo = pExprs;
310,523,801✔
2833
  return code;
310,505,833✔
2834
}
2835

2836
static void deleteSubsidiareCtx(void* pData) {
×
2837
  SSubsidiaryResInfo* pCtx = (SSubsidiaryResInfo*)pData;
×
2838
  if (pCtx->pCtx) {
×
2839
    taosMemoryFreeClear(pCtx->pCtx);
×
2840
  }
2841
}
×
2842

2843
// set the output buffer for the selectivity + tag query
2844
static int32_t setSelectValueColumnInfo(SqlFunctionCtx* pCtx, int32_t numOfOutput) {
332,798,892✔
2845
  int32_t num = 0;
332,798,892✔
2846
  int32_t code = TSDB_CODE_SUCCESS;
332,798,892✔
2847
  int32_t lino = 0;
332,798,892✔
2848

2849
  SArray* pValCtxArray = NULL;
332,798,892✔
2850
  for (int32_t i = numOfOutput - 1; i > 0; --i) {  // select Func is at the end of the list
831,608,927✔
2851
    int32_t funcIdx = pCtx[i].pExpr->pExpr->_function.bindExprID;
498,864,965✔
2852
    if (funcIdx > 0) {
498,876,465✔
2853
      if (pValCtxArray == NULL) {
1,838,806✔
2854
        // the end of the list is the select function of biggest index
2855
        pValCtxArray = taosArrayInit_s(sizeof(SSubsidiaryResInfo*), funcIdx);
1,321,938✔
2856
        if (pValCtxArray == NULL) {
1,319,428✔
2857
          return terrno;
×
2858
        }
2859
      }
2860
      if (funcIdx > pValCtxArray->size) {
1,836,296✔
2861
        qError("funcIdx:%d is out of range", funcIdx);
×
2862
        taosArrayDestroyP(pValCtxArray, deleteSubsidiareCtx);
×
2863
        return TSDB_CODE_QRY_EXECUTOR_INTERNAL_ERROR;
×
2864
      }
2865
      SSubsidiaryResInfo* pSubsidiary = &pCtx[i].subsidiaries;
1,836,296✔
2866
      pSubsidiary->pCtx = taosMemoryCalloc(numOfOutput, POINTER_BYTES);
1,836,296✔
2867
      if (pSubsidiary->pCtx == NULL) {
1,838,304✔
2868
        taosArrayDestroyP(pValCtxArray, deleteSubsidiareCtx);
×
2869
        return terrno;
×
2870
      }
2871
      pSubsidiary->num = 0;
1,834,288✔
2872
      taosArraySet(pValCtxArray, funcIdx - 1, &pSubsidiary);
1,835,292✔
2873
    }
2874
  }
2875

2876
  SqlFunctionCtx*  p = NULL;
332,743,962✔
2877
  SqlFunctionCtx** pValCtx = NULL;
332,743,962✔
2878
  if (pValCtxArray == NULL) {
332,743,962✔
2879
    pValCtx = taosMemoryCalloc(numOfOutput, POINTER_BYTES);
331,463,448✔
2880
    if (pValCtx == NULL) {
331,382,535✔
2881
      QUERY_CHECK_CODE(terrno, lino, _end);
×
2882
    }
2883
  }
2884

2885
  for (int32_t i = 0; i < numOfOutput; ++i) {
1,132,260,500✔
2886
    const char* pName = pCtx[i].pExpr->pExpr->_function.functionName;
799,647,400✔
2887
    if ((strcmp(pName, "_select_value") == 0)) {
799,777,789✔
2888
      if (pValCtxArray == NULL) {
5,984,149✔
2889
        pValCtx[num++] = &pCtx[i];
3,405,154✔
2890
      } else {
2891
        int32_t bindFuncIndex = pCtx[i].pExpr->pExpr->relatedTo;  // start from index 1;
2,578,995✔
2892
        if (bindFuncIndex > 0) {                                  // 0 is default index related to the select function
2,582,007✔
2893
          bindFuncIndex -= 1;
2,516,747✔
2894
        }
2895
        SSubsidiaryResInfo** pSubsidiary = taosArrayGet(pValCtxArray, bindFuncIndex);
2,582,007✔
2896
        if (pSubsidiary == NULL) {
2,578,995✔
2897
          QUERY_CHECK_CODE(TSDB_CODE_QRY_EXECUTOR_INTERNAL_ERROR, lino, _end);
×
2898
        }
2899
        (*pSubsidiary)->pCtx[(*pSubsidiary)->num] = &pCtx[i];
2,578,995✔
2900
        (*pSubsidiary)->num++;
2,573,975✔
2901
      }
2902
    } else if (fmIsSelectFunc(pCtx[i].functionId)) {
793,793,640✔
2903
      if (pValCtxArray == NULL) {
65,913,343✔
2904
        p = &pCtx[i];
63,669,485✔
2905
      }
2906
    }
2907
  }
2908

2909
  if (p != NULL) {
332,613,100✔
2910
    p->subsidiaries.pCtx = pValCtx;
25,572,224✔
2911
    p->subsidiaries.num = num;
25,572,721✔
2912
  } else {
2913
    taosMemoryFreeClear(pValCtx);
307,040,876✔
2914
  }
2915

2916
_end:
1,276,180✔
2917
  if (code != TSDB_CODE_SUCCESS) {
332,630,727✔
2918
    taosArrayDestroyP(pValCtxArray, deleteSubsidiareCtx);
×
2919
    taosMemoryFreeClear(pValCtx);
×
2920
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
2921
  } else {
2922
    taosArrayDestroy(pValCtxArray);
332,630,727✔
2923
  }
2924
  return code;
332,714,267✔
2925
}
2926

2927
SqlFunctionCtx* createSqlFunctionCtx(SExprInfo* pExprInfo, int32_t numOfOutput, int32_t** rowEntryInfoOffset,
332,826,749✔
2928
                                     SFunctionStateStore* pStore) {
2929
  int32_t         code = TSDB_CODE_SUCCESS;
332,826,749✔
2930
  int32_t         lino = 0;
332,826,749✔
2931
  SqlFunctionCtx* pFuncCtx = (SqlFunctionCtx*)taosMemoryCalloc(numOfOutput, sizeof(SqlFunctionCtx));
332,826,749✔
2932
  if (pFuncCtx == NULL) {
332,510,100✔
2933
    return NULL;
×
2934
  }
2935

2936
  *rowEntryInfoOffset = taosMemoryCalloc(numOfOutput, sizeof(int32_t));
332,510,100✔
2937
  if (*rowEntryInfoOffset == 0) {
332,743,266✔
2938
    taosMemoryFreeClear(pFuncCtx);
×
2939
    return NULL;
×
2940
  }
2941

2942
  for (int32_t i = 0; i < numOfOutput; ++i) {
1,132,481,071✔
2943
    SExprInfo* pExpr = &pExprInfo[i];
799,813,987✔
2944

2945
    SExprBasicInfo* pFunct = &pExpr->base;
799,713,090✔
2946
    SqlFunctionCtx* pCtx = &pFuncCtx[i];
799,757,987✔
2947

2948
    pCtx->functionId = -1;
799,766,555✔
2949
    pCtx->pExpr = pExpr;
799,813,733✔
2950

2951
    if (pExpr->pExpr->nodeType == QUERY_NODE_FUNCTION) {
799,721,176✔
2952
      SFuncExecEnv env = {0};
254,087,706✔
2953
      pCtx->functionId = pExpr->pExpr->_function.pFunctNode->funcId;
254,081,891✔
2954
      pCtx->isPseudoFunc = fmIsWindowPseudoColumnFunc(pCtx->functionId) || fmIsPlaceHolderFunc(pCtx->functionId);
254,088,356✔
2955
      pCtx->isNotNullFunc = fmIsNotNullOutputFunc(pCtx->functionId);
254,096,183✔
2956

2957
      bool isUdaf = fmIsUserDefinedFunc(pCtx->functionId);
254,077,070✔
2958
      if (fmIsAggFunc(pCtx->functionId) || fmIsIndefiniteRowsFunc(pCtx->functionId)) {
421,409,887✔
2959
        if (!isUdaf) {
167,417,297✔
2960
          code = fmGetFuncExecFuncs(pCtx->functionId, &pCtx->fpSet);
167,374,975✔
2961
          QUERY_CHECK_CODE(code, lino, _end);
167,357,862✔
2962
        } else {
2963
          char* udfName = pExpr->pExpr->_function.pFunctNode->functionName;
42,322✔
2964
          pCtx->udfName = taosStrdup(udfName);
42,322✔
2965
          QUERY_CHECK_NULL(pCtx->udfName, code, lino, _end, terrno);
42,322✔
2966

2967
          code = fmGetUdafExecFuncs(pCtx->functionId, &pCtx->fpSet);
42,322✔
2968
          QUERY_CHECK_CODE(code, lino, _end);
42,322✔
2969
        }
2970
        bool tmp = pCtx->fpSet.getEnv(pExpr->pExpr->_function.pFunctNode, &env);
167,400,184✔
2971
        if (!tmp) {
167,370,547✔
2972
          code = terrno;
×
2973
          QUERY_CHECK_CODE(code, lino, _end);
4,153✔
2974
        }
2975
      } else {
2976
        if (fmIsPlaceHolderFunc(pCtx->functionId)) {
86,636,985✔
2977
          code = fmGetStreamPesudoFuncEnv(pCtx->functionId, pExpr->base.pParamList, &env);
8,184,647✔
2978
          QUERY_CHECK_CODE(code, lino, _end);
8,184,647✔
2979
        }      
2980
        
2981
        code = fmGetScalarFuncExecFuncs(pCtx->functionId, &pCtx->sfp);
86,666,226✔
2982
        if (code != TSDB_CODE_SUCCESS && isUdaf) {
86,642,104✔
2983
          code = TSDB_CODE_SUCCESS;
25,992✔
2984
        }
2985
        QUERY_CHECK_CODE(code, lino, _end);
86,642,104✔
2986

2987
        if (pCtx->sfp.getEnv != NULL) {
86,642,104✔
2988
          bool tmp = pCtx->sfp.getEnv(pExpr->pExpr->_function.pFunctNode, &env);
16,991,288✔
2989
          if (!tmp) {
16,994,544✔
2990
            code = terrno;
×
2991
            QUERY_CHECK_CODE(code, lino, _end);
×
2992
          }
2993
        }
2994
      }
2995
      pCtx->resDataInfo.interBufSize = env.calcMemSize;
254,030,711✔
2996
    } else if (pExpr->pExpr->nodeType == QUERY_NODE_COLUMN || pExpr->pExpr->nodeType == QUERY_NODE_OPERATOR ||
545,629,163✔
2997
               pExpr->pExpr->nodeType == QUERY_NODE_VALUE) {
39,477,094✔
2998
      // for simple column, the result buffer needs to hold at least one element.
2999
      pCtx->resDataInfo.interBufSize = pFunct->resSchema.bytes;
545,814,040✔
3000
    }
3001

3002
    pCtx->input.numOfInputCols = pFunct->numOfParams;
799,869,487✔
3003
    pCtx->input.pData = taosMemoryCalloc(pFunct->numOfParams, POINTER_BYTES);
799,782,010✔
3004
    QUERY_CHECK_NULL(pCtx->input.pData, code, lino, _end, terrno);
799,848,561✔
3005
    pCtx->input.pColumnDataAgg = taosMemoryCalloc(pFunct->numOfParams, POINTER_BYTES);
799,763,650✔
3006
    QUERY_CHECK_NULL(pCtx->input.pColumnDataAgg, code, lino, _end, terrno);
799,831,272✔
3007

3008
    pCtx->pTsOutput = NULL;
799,763,640✔
3009
    pCtx->resDataInfo.bytes = pFunct->resSchema.bytes;
799,898,222✔
3010
    pCtx->resDataInfo.type = pFunct->resSchema.type;
799,819,309✔
3011
    pCtx->order = TSDB_ORDER_ASC;
799,875,427✔
3012
    pCtx->start.key = INT64_MIN;
799,968,562✔
3013
    pCtx->end.key = INT64_MIN;
799,816,599✔
3014
    pCtx->numOfParams = pExpr->base.numOfParams;
799,777,095✔
3015
    pCtx->param = pFunct->pParam;
800,028,672✔
3016
    pCtx->saveHandle.currentPage = -1;
799,863,683✔
3017
    pCtx->pStore = pStore;
799,845,959✔
3018
    pCtx->hasWindowOrGroup = false;
799,917,021✔
3019
    pCtx->needCleanup = false;
799,813,293✔
3020
    pCtx->skipDynDataCheck = false;
799,697,098✔
3021
  }
3022

3023
  for (int32_t i = 1; i < numOfOutput; ++i) {
831,726,195✔
3024
    (*rowEntryInfoOffset)[i] = (int32_t)((*rowEntryInfoOffset)[i - 1] + sizeof(SResultRowEntryInfo) +
997,846,691✔
3025
                                         pFuncCtx[i - 1].resDataInfo.interBufSize);
498,954,082✔
3026
  }
3027

3028
  code = setSelectValueColumnInfo(pFuncCtx, numOfOutput);
332,786,764✔
3029
  QUERY_CHECK_CODE(code, lino, _end);
332,718,432✔
3030

3031
_end:
332,718,432✔
3032
  if (code != TSDB_CODE_SUCCESS) {
332,667,379✔
3033
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
3034
    for (int32_t i = 0; i < numOfOutput; ++i) {
×
3035
      taosMemoryFree(pFuncCtx[i].input.pData);
×
3036
      taosMemoryFree(pFuncCtx[i].input.pColumnDataAgg);
×
3037
    }
3038
    taosMemoryFreeClear(*rowEntryInfoOffset);
×
3039
    taosMemoryFreeClear(pFuncCtx);
×
3040

3041
    terrno = code;
×
3042
    return NULL;
×
3043
  }
3044
  return pFuncCtx;
332,667,379✔
3045
}
3046

3047
// NOTE: sources columns are more than the destination SSDatablock columns.
3048
// doFilter in table scan needs every column even its output is false
3049
int32_t relocateColumnData(SSDataBlock* pBlock, const SArray* pColMatchInfo, SArray* pCols, bool outputEveryColumn) {
10,073,258✔
3050
  int32_t code = TSDB_CODE_SUCCESS;
10,073,258✔
3051
  size_t  numOfSrcCols = taosArrayGetSize(pCols);
10,073,258✔
3052

3053
  int32_t i = 0, j = 0;
10,073,258✔
3054
  while (i < numOfSrcCols && j < taosArrayGetSize(pColMatchInfo)) {
93,304,715✔
3055
    SColumnInfoData* p = taosArrayGet(pCols, i);
83,231,457✔
3056
    if (!p) {
83,230,817✔
3057
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
3058
      return terrno;
×
3059
    }
3060
    SColMatchItem* pmInfo = taosArrayGet(pColMatchInfo, j);
83,230,817✔
3061
    if (!pmInfo) {
83,231,086✔
3062
      return terrno;
×
3063
    }
3064

3065
    if (p->info.colId == pmInfo->colId) {
83,231,086✔
3066
      SColumnInfoData* pDst = taosArrayGet(pBlock->pDataBlock, pmInfo->dstSlotId);
74,846,556✔
3067
      if (!pDst) {
74,845,366✔
3068
        return terrno;
×
3069
      }
3070
      code = colDataAssign(pDst, p, pBlock->info.rows, &pBlock->info);
74,845,366✔
3071
      if (code != TSDB_CODE_SUCCESS) {
74,846,376✔
3072
        qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(code));
×
3073
        return code;
×
3074
      }
3075
      i++;
74,846,376✔
3076
      j++;
74,846,376✔
3077
    } else if (p->info.colId < pmInfo->colId) {
8,385,081✔
3078
      i++;
8,385,081✔
3079
    } else {
3080
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(TSDB_CODE_QRY_EXECUTOR_INTERNAL_ERROR));
×
3081
      return TSDB_CODE_QRY_EXECUTOR_INTERNAL_ERROR;
×
3082
    }
3083
  }
3084
  return code;
10,073,258✔
3085
}
3086

3087
SInterval extractIntervalInfo(const STableScanPhysiNode* pTableScanNode) {
193,685,405✔
3088
  SInterval interval = {
387,126,290✔
3089
      .interval = pTableScanNode->interval,
193,447,112✔
3090
      .sliding = pTableScanNode->sliding,
193,438,618✔
3091
      .intervalUnit = pTableScanNode->intervalUnit,
193,746,527✔
3092
      .slidingUnit = pTableScanNode->slidingUnit,
193,426,019✔
3093
      .offset = pTableScanNode->offset,
193,751,865✔
3094
      .precision = pTableScanNode->scan.node.pOutputDataBlockDesc->precision,
193,745,330✔
3095
      .timeRange = pTableScanNode->scanRange,
3096
  };
3097
  calcIntervalAutoOffset(&interval);
193,444,062✔
3098

3099
  return interval;
193,511,144✔
3100
}
3101

3102
SColumn extractColumnFromColumnNode(SColumnNode* pColNode) {
45,875,944✔
3103
  SColumn c = {0};
45,875,944✔
3104

3105
  c.slotId = pColNode->slotId;
45,875,944✔
3106
  c.colId = pColNode->colId;
45,879,211✔
3107
  c.type = pColNode->node.resType.type;
45,882,764✔
3108
  c.bytes = pColNode->node.resType.bytes;
45,873,308✔
3109
  c.scale = pColNode->node.resType.scale;
45,876,552✔
3110
  c.precision = pColNode->node.resType.precision;
45,881,375✔
3111
  return c;
45,862,650✔
3112
}
3113

3114

3115
/**
3116
 * @brief Determine the actual time range for reading data based on the RANGE clause and the WHERE conditions.
3117
 * @param[in] cond The range specified by WHERE condition.
3118
 * @param[in] range The range specified by RANGE clause.
3119
 * @param[out] twindow The range to be read in DESC order, and only one record is needed.
3120
 * @param[out] extTwindow The external range to read for only one record, which is used for FILL clause.
3121
 * @note `cond` and `twindow` may be the same address.
3122
 */
3123
static int32_t getQueryExtWindow(const STimeWindow* cond, const STimeWindow* range, STimeWindow* twindow,
1,643,241✔
3124
                                 STimeWindow* extTwindows) {
3125
  int32_t     code = TSDB_CODE_SUCCESS;
1,643,241✔
3126
  int32_t     lino = 0;
1,643,241✔
3127
  STimeWindow tempWindow;
3128

3129
  if (cond->skey > cond->ekey || range->skey > range->ekey) {
1,643,241✔
3130
    *twindow = extTwindows[0] = extTwindows[1] = TSWINDOW_DESC_INITIALIZER;
3,005✔
3131
    return code;
3,005✔
3132
  }
3133

3134
  if (range->ekey < cond->skey) {
1,640,236✔
3135
    extTwindows[1] = *cond;
261,032✔
3136
    *twindow = extTwindows[0] = TSWINDOW_DESC_INITIALIZER;
261,032✔
3137
    return code;
261,032✔
3138
  }
3139

3140
  if (cond->ekey < range->skey) {
1,379,204✔
3141
    extTwindows[0] = *cond;
193,338✔
3142
    *twindow = extTwindows[1] = TSWINDOW_DESC_INITIALIZER;
193,338✔
3143
    return code;
193,338✔
3144
  }
3145

3146
  // Only scan data in the time range intersecion.
3147
  extTwindows[0] = extTwindows[1] = *cond;
1,185,866✔
3148
  twindow->skey = TMAX(cond->skey, range->skey);
1,185,866✔
3149
  twindow->ekey = TMIN(cond->ekey, range->ekey);
1,185,866✔
3150
  extTwindows[0].ekey = twindow->skey - 1;
1,185,866✔
3151
  extTwindows[1].skey = twindow->ekey + 1;
1,185,866✔
3152

3153
  return code;
1,185,866✔
3154
}
3155

3156
static int32_t getPrimaryTimeRange(SNode** pPrimaryKeyCond, STimeWindow* pTimeRange, bool* isStrict) {
12,360✔
3157
  SNode*  pNew = NULL;
12,360✔
3158
  int32_t code = scalarCalculateRemoteConstants(*pPrimaryKeyCond, &pNew);
12,360✔
3159
  if (TSDB_CODE_SUCCESS == code) {
12,360✔
3160
    *pPrimaryKeyCond = pNew;
12,360✔
3161
    if (nodeType(pNew) != QUERY_NODE_VALUE) {
12,360✔
3162
      code = filterGetTimeRange(*pPrimaryKeyCond, pTimeRange, isStrict, NULL);
12,360✔
3163
    }
3164
  }
3165
  return code;
12,360✔
3166
}
3167

3168
int32_t initQueryTableDataCond(SQueryTableDataCond* pCond, STableScanPhysiNode* pTableScanNode,
215,993,594✔
3169
                               const SReadHandle* readHandle, bool applyExtWin) {
3170
  int32_t code = 0;                             
215,993,594✔
3171
  pCond->order = pTableScanNode->scanSeq[0] > 0 ? TSDB_ORDER_ASC : TSDB_ORDER_DESC;
215,993,594✔
3172
  pCond->numOfCols = LIST_LENGTH(pTableScanNode->scan.pScanCols);
215,984,976✔
3173

3174
  pCond->colList = taosMemoryCalloc(pCond->numOfCols, sizeof(SColumnInfo));
216,045,510✔
3175
  if (!pCond->colList) {
215,888,478✔
3176
    return terrno;
×
3177
  }
3178
  pCond->pSlotList = taosMemoryMalloc(sizeof(int32_t) * pCond->numOfCols);
215,805,433✔
3179
  if (pCond->pSlotList == NULL) {
215,931,393✔
3180
    taosMemoryFreeClear(pCond->colList);
×
3181
    return terrno;
×
3182
  }
3183

3184
  // TODO: get it from stable scan node
3185
  pCond->twindows = pTableScanNode->scanRange;
215,791,099✔
3186
  pCond->suid = pTableScanNode->scan.suid;
216,038,837✔
3187
  pCond->type = TIMEWINDOW_RANGE_CONTAINED;
215,838,977✔
3188
  pCond->startVersion = -1;
215,960,075✔
3189
  pCond->endVersion = -1;
216,031,781✔
3190
  pCond->skipRollup = readHandle->skipRollup;
215,688,188✔
3191
  if (readHandle->winRangeValid) {
215,895,929✔
3192
    pCond->twindows = readHandle->winRange;
347,914✔
3193
  }
3194
  pCond->cacheSttStatis = readHandle->cacheSttStatis;
216,033,602✔
3195
  // allowed read stt file optimization mode
3196
  pCond->notLoadData = (pTableScanNode->dataRequired == FUNC_DATA_REQUIRED_NOT_LOAD) &&
432,022,087✔
3197
                       (pTableScanNode->scan.node.pConditions == NULL) && (pTableScanNode->interval == 0);
215,957,929✔
3198

3199
  int32_t j = 0;
215,928,120✔
3200
  for (int32_t i = 0; i < pCond->numOfCols; ++i) {
1,023,899,763✔
3201
    STargetNode* pNode = (STargetNode*)nodesListGetNode(pTableScanNode->scan.pScanCols, i);
808,148,971✔
3202
    if (!pNode) {
807,643,831✔
3203
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
3204
      return terrno;
×
3205
    }
3206
    SColumnNode* pColNode = (SColumnNode*)pNode->pExpr;
807,643,831✔
3207
    if (pColNode->colType == COLUMN_TYPE_TAG) {
808,065,520✔
3208
      continue;
×
3209
    }
3210

3211
    pCond->colList[j].type = pColNode->node.resType.type;
807,987,625✔
3212
    pCond->colList[j].bytes = pColNode->node.resType.bytes;
807,995,632✔
3213
    pCond->colList[j].colId = pColNode->colId;
807,880,370✔
3214
    pCond->colList[j].pk = pColNode->isPk;
808,065,562✔
3215

3216
    pCond->pSlotList[j] = pNode->slotId;
808,159,946✔
3217
    j += 1;
807,971,643✔
3218
  }
3219

3220
  pCond->numOfCols = j;
216,089,552✔
3221

3222
  if (applyExtWin) {
216,089,174✔
3223
    if (NULL != pTableScanNode->pExtScanRange) {
194,159,298✔
3224
      pCond->type = TIMEWINDOW_RANGE_EXTERNAL;
1,643,241✔
3225
      code = getQueryExtWindow(&pCond->twindows, pTableScanNode->pExtScanRange, &pCond->twindows, pCond->extTwindows);
1,643,241✔
3226
    } else if (readHandle->extWinRangeValid) {
192,272,364✔
3227
      pCond->type = TIMEWINDOW_RANGE_EXTERNAL;
×
3228
      code = getQueryExtWindow(&pCond->twindows, &readHandle->extWinRange, &pCond->twindows, pCond->extTwindows);
×
3229
    }
3230
  }
3231

3232
  if (pTableScanNode->pPrimaryCond) {
215,919,824✔
3233
    bool isStrict = false;
12,360✔
3234
    code = getPrimaryTimeRange((SNode**)&pTableScanNode->pPrimaryCond, &pCond->twindows, &isStrict);
12,360✔
3235
    if (code || !isStrict) {
12,360✔
3236
      code = nodesMergeNode((SNode**)&pTableScanNode->scan.node.pConditions, &pTableScanNode->pPrimaryCond);
×
3237
    }
3238
  }
3239

3240
  return code;
215,973,661✔
3241
}
3242

3243
int32_t initQueryTableDataCondWithColArray(SQueryTableDataCond* pCond, SQueryTableDataCond* pOrgCond,
15,003,514✔
3244
                                           const SReadHandle* readHandle, SArray* colArray) {
3245
  int32_t code = TSDB_CODE_SUCCESS;
15,003,514✔
3246
  int32_t lino = 0;
15,003,514✔
3247

3248
  pCond->order = TSDB_ORDER_ASC;
15,003,514✔
3249
  pCond->numOfCols = (int32_t)taosArrayGetSize(colArray);
15,005,164✔
3250

3251
  pCond->colList = taosMemoryCalloc(pCond->numOfCols, sizeof(SColumnInfo));
15,000,214✔
3252
  QUERY_CHECK_NULL(pCond->colList, code, lino, _return, terrno);
15,001,864✔
3253

3254
  pCond->pSlotList = taosMemoryMalloc(sizeof(int32_t) * pCond->numOfCols);
14,999,664✔
3255
  QUERY_CHECK_NULL(pCond->pSlotList, code, lino, _return, terrno);
15,000,764✔
3256

3257
  pCond->twindows = pOrgCond->twindows;
14,990,314✔
3258
  pCond->order = pOrgCond->order;
15,000,214✔
3259
  pCond->type = pOrgCond->type;
15,000,764✔
3260
  pCond->startVersion = -1;
15,000,214✔
3261
  pCond->endVersion = -1;
14,996,914✔
3262
  pCond->skipRollup = true;
15,002,964✔
3263
  pCond->notLoadData = false;
14,996,914✔
3264

3265
  for (int32_t i = 0; i < pCond->numOfCols; ++i) {
72,538,128✔
3266
    SColIdPair* pColPair = taosArrayGet(colArray, i);
57,541,764✔
3267
    QUERY_CHECK_NULL(pColPair, code, lino, _return, terrno);
57,545,064✔
3268

3269
    bool find = false;
57,542,314✔
3270
    for (int32_t j = 0; j < pOrgCond->numOfCols; ++j) {
344,886,676✔
3271
      if (pOrgCond->colList[j].colId == pColPair->vtbColId) {
344,771,726✔
3272
        pCond->colList[i].type = pOrgCond->colList[j].type;
57,562,664✔
3273
        pCond->colList[i].bytes = pOrgCond->colList[j].bytes;
57,562,114✔
3274
        pCond->colList[i].colId = pColPair->orgColId;
57,563,764✔
3275
        pCond->colList[i].pk = pOrgCond->colList[j].pk;
57,563,214✔
3276
        pCond->pSlotList[i] = i;
57,564,864✔
3277
        find = true;
57,556,064✔
3278
        qDebug("%s mapped vtb colId:%d to org colId:%d", __func__, pColPair->vtbColId, pColPair->orgColId);
57,556,064✔
3279
        break;
57,544,514✔
3280
      }
3281
    }
3282
    QUERY_CHECK_CONDITION(find, code, lino, _return, TSDB_CODE_NOT_FOUND);
57,542,864✔
3283
  }
3284

3285
  return code;
15,007,914✔
3286
_return:
×
3287
  qError("%s failed at line %d since %s", __func__, lino, tstrerror(terrno));
×
3288
  taosMemoryFreeClear(pCond->colList);
×
3289
  taosMemoryFreeClear(pCond->pSlotList);
×
3290
  return code;
×
3291
}
3292

3293
void cleanupQueryTableDataCond(SQueryTableDataCond* pCond) {
469,844,608✔
3294
  taosMemoryFreeClear(pCond->colList);
469,844,608✔
3295
  taosMemoryFreeClear(pCond->pSlotList);
469,755,138✔
3296
}
469,718,733✔
3297

3298
int32_t convertFillType(int32_t mode) {
2,077,246✔
3299
  int32_t type = TSDB_FILL_NONE;
2,077,246✔
3300
  switch (mode) {
2,077,246✔
3301
    case FILL_MODE_PREV:
113,210✔
3302
      type = TSDB_FILL_PREV;
113,210✔
3303
      break;
113,210✔
3304
    case FILL_MODE_NONE:
×
3305
      type = TSDB_FILL_NONE;
×
3306
      break;
×
3307
    case FILL_MODE_NULL:
133,858✔
3308
      type = TSDB_FILL_NULL;
133,858✔
3309
      break;
133,858✔
3310
    case FILL_MODE_NULL_F:
15,841✔
3311
      type = TSDB_FILL_NULL_F;
15,841✔
3312
      break;
15,841✔
3313
    case FILL_MODE_NEXT:
98,408✔
3314
      type = TSDB_FILL_NEXT;
98,408✔
3315
      break;
98,408✔
3316
    case FILL_MODE_VALUE:
150,206✔
3317
      type = TSDB_FILL_SET_VALUE;
150,206✔
3318
      break;
150,206✔
3319
    case FILL_MODE_VALUE_F:
4,324✔
3320
      type = TSDB_FILL_SET_VALUE_F;
4,324✔
3321
      break;
4,324✔
3322
    case FILL_MODE_LINEAR:
146,791✔
3323
      type = TSDB_FILL_LINEAR;
146,791✔
3324
      break;
146,791✔
3325
    case FILL_MODE_NEAR:
1,414,113✔
3326
      type = TSDB_FILL_NEAR;
1,414,113✔
3327
      break;
1,414,113✔
3328
    default:
495✔
3329
      type = TSDB_FILL_NONE;
495✔
3330
  }
3331

3332
  return type;
2,077,246✔
3333
}
3334

3335
void getInitialStartTimeWindow(SInterval* pInterval, TSKEY ts, STimeWindow* w, bool ascQuery) {
1,793,488,817✔
3336
  if (ascQuery) {
1,793,488,817✔
3337
    *w = getAlignQueryTimeWindow(pInterval, ts);
1,794,286,712✔
3338
  } else {
3339
    // the start position of the first time window in the endpoint that spreads beyond the queried last timestamp
3340
    *w = getAlignQueryTimeWindow(pInterval, ts);
4,174✔
3341

3342
    int64_t key = w->skey;
146,589✔
3343
    while (key < ts) {  // moving towards end
161,364✔
3344
      key = getNextTimeWindowStart(pInterval, key, TSDB_ORDER_ASC);
76,853✔
3345
      if (key > ts) {
76,853✔
3346
        break;
62,078✔
3347
      }
3348

3349
      w->skey = key;
14,775✔
3350
    }
3351
    w->ekey = taosTimeAdd(w->skey, pInterval->interval, pInterval->intervalUnit, pInterval->precision, NULL) - 1;
146,589✔
3352
  }
3353
}
1,795,069,333✔
3354

3355
static STimeWindow doCalculateTimeWindow(int64_t ts, SInterval* pInterval) {
26,634,286✔
3356
  STimeWindow w = {0};
26,634,286✔
3357

3358
  w.skey = taosTimeTruncate(ts, pInterval);
26,634,286✔
3359
  w.ekey = taosTimeGetIntervalEnd(w.skey, pInterval);
26,633,640✔
3360
  return w;
26,634,550✔
3361
}
3362

3363
STimeWindow getFirstQualifiedTimeWindow(int64_t ts, STimeWindow* pWindow, SInterval* pInterval, int32_t order) {
1,083,989✔
3364
  STimeWindow win = *pWindow;
1,083,989✔
3365
  STimeWindow save = win;
1,083,989✔
3366
  while (win.skey <= ts && win.ekey >= ts) {
5,482,324✔
3367
    save = win;
4,398,335✔
3368
    // get previous time window
3369
    getNextTimeWindow(pInterval, &win, order == TSDB_ORDER_DESC ? TSDB_ORDER_ASC : TSDB_ORDER_DESC);
4,398,335✔
3370
  }
3371

3372
  return save;
1,083,989✔
3373
}
3374

3375
// get the correct time window according to the handled timestamp
3376
// todo refactor
3377
STimeWindow getActiveTimeWindow(SDiskbasedBuf* pBuf, SResultRowInfo* pResultRowInfo, int64_t ts, SInterval* pInterval,
44,461,530✔
3378
                                int32_t order) {
3379
  STimeWindow w = {0};
44,461,530✔
3380
  if (pResultRowInfo->cur.pageId == -1) {  // the first window, from the previous stored value
44,460,376✔
3381
    getInitialStartTimeWindow(pInterval, ts, &w, (order != TSDB_ORDER_DESC));
5,984,143✔
3382
    return w;
5,989,163✔
3383
  }
3384

3385
  SResultRow* pRow = getResultRowByPos(pBuf, &pResultRowInfo->cur, false);
38,472,274✔
3386
  if (pRow) {
38,471,208✔
3387
    TAOS_SET_OBJ_ALIGNED(&w, pRow->win);
38,471,355✔
3388
  }
3389

3390
  // in case of typical time window, we can calculate time window directly.
3391
  if (w.skey > ts || w.ekey < ts) {
38,471,446✔
3392
    w = doCalculateTimeWindow(ts, pInterval);
26,633,946✔
3393
  }
3394

3395
  if (pInterval->interval != pInterval->sliding) {
38,472,155✔
3396
    // it is an sliding window query, in which sliding value is not equalled to
3397
    // interval value, and we need to find the first qualified time window.
3398
    w = getFirstQualifiedTimeWindow(ts, &w, pInterval, order);
1,083,989✔
3399
  }
3400

3401
  return w;
38,470,669✔
3402
}
3403

3404
TSKEY getNextTimeWindowStart(const SInterval* pInterval, TSKEY start, int32_t order) {
2,147,483,647✔
3405
  int32_t factor = GET_FORWARD_DIRECTION_FACTOR(order);
2,147,483,647✔
3406
  TSKEY   nextStart = taosTimeAdd(start, -1 * pInterval->offset, pInterval->offsetUnit, pInterval->precision, NULL);
2,147,483,647✔
3407
  nextStart = taosTimeAdd(nextStart, factor * pInterval->sliding, pInterval->slidingUnit, pInterval->precision, NULL);
2,147,483,647✔
3408
  nextStart = taosTimeAdd(nextStart, pInterval->offset, pInterval->offsetUnit, pInterval->precision, NULL);
2,147,483,647✔
3409
  return nextStart;
2,147,483,647✔
3410
}
3411

3412
void getNextTimeWindow(const SInterval* pInterval, STimeWindow* tw, int32_t order) {
2,147,483,647✔
3413
  tw->skey = getNextTimeWindowStart(pInterval, tw->skey, order);
2,147,483,647✔
3414
  tw->ekey = taosTimeAdd(tw->skey, pInterval->interval, pInterval->intervalUnit, pInterval->precision, NULL) - 1;
2,147,483,647✔
3415
}
2,147,483,647✔
3416

3417
bool hasLimitOffsetInfo(SLimitInfo* pLimitInfo) {
314,853,848✔
3418
  return (pLimitInfo->limit.limit != -1 || pLimitInfo->limit.offset != -1 || pLimitInfo->slimit.limit != -1 ||
626,927,512✔
3419
          pLimitInfo->slimit.offset != -1);
312,073,862✔
3420
}
3421

3422
bool hasSlimitOffsetInfo(SLimitInfo* pLimitInfo) {
×
3423
  return (pLimitInfo->slimit.limit != -1 || pLimitInfo->slimit.offset != -1);
×
3424
}
3425

3426
void initLimitInfo(const SNode* pLimit, const SNode* pSLimit, SLimitInfo* pLimitInfo) {
474,696,247✔
3427
  SLimit limit = {.limit = getLimit(pLimit), .offset = getOffset(pLimit)};
474,696,247✔
3428
  SLimit slimit = {.limit = getLimit(pSLimit), .offset = getOffset(pSLimit)};
474,448,534✔
3429

3430
  pLimitInfo->limit = limit;
474,409,761✔
3431
  pLimitInfo->slimit = slimit;
474,472,304✔
3432
  pLimitInfo->remainOffset = limit.offset;
474,495,475✔
3433
  pLimitInfo->remainGroupOffset = slimit.offset;
474,491,932✔
3434
  pLimitInfo->numOfOutputRows = 0;
474,471,045✔
3435
  pLimitInfo->numOfOutputGroups = 0;
474,654,774✔
3436
  pLimitInfo->currentGroupId = 0;
474,634,006✔
3437
}
474,689,023✔
3438

3439
void resetLimitInfoForNextGroup(SLimitInfo* pLimitInfo) {
55,017,661✔
3440
  pLimitInfo->numOfOutputRows = 0;
55,017,661✔
3441
  pLimitInfo->remainOffset = pLimitInfo->limit.offset;
55,037,402✔
3442
}
55,008,688✔
3443

3444
int32_t tableListGetSize(const STableListInfo* pTableList, int32_t* pRes) {
487,708,766✔
3445
  if (taosArrayGetSize(pTableList->pTableList) != taosHashGetSize(pTableList->map)) {
487,708,766✔
3446
    qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(TSDB_CODE_QRY_EXECUTOR_INTERNAL_ERROR));
×
3447
    return TSDB_CODE_QRY_EXECUTOR_INTERNAL_ERROR;
×
3448
  }
3449
  (*pRes) = taosArrayGetSize(pTableList->pTableList);
487,640,172✔
3450
  return TSDB_CODE_SUCCESS;
487,549,879✔
3451
}
3452

3453
uint64_t tableListGetSuid(const STableListInfo* pTableList) { return pTableList->idInfo.suid; }
3,237,206✔
3454

3455
STableKeyInfo* tableListGetInfo(const STableListInfo* pTableList, int32_t index) {
155,989,262✔
3456
  if (taosArrayGetSize(pTableList->pTableList) == 0) {
155,989,262✔
3457
    return NULL;
3,675✔
3458
  }
3459

3460
  return taosArrayGet(pTableList->pTableList, index);
155,973,152✔
3461
}
3462

3463
int32_t tableListFind(const STableListInfo* pTableList, uint64_t uid, int32_t startIndex) {
9,676✔
3464
  int32_t numOfTables = taosArrayGetSize(pTableList->pTableList);
9,676✔
3465
  if (startIndex >= numOfTables) {
9,676✔
3466
    return -1;
×
3467
  }
3468

3469
  for (int32_t i = startIndex; i < numOfTables; ++i) {
118,906✔
3470
    STableKeyInfo* p = taosArrayGet(pTableList->pTableList, i);
118,906✔
3471
    if (!p) {
118,906✔
3472
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
3473
      return -1;
×
3474
    }
3475
    if (p->uid == uid) {
118,906✔
3476
      return i;
9,676✔
3477
    }
3478
  }
3479
  return -1;
×
3480
}
3481

3482
void tableListGetSourceTableInfo(const STableListInfo* pTableList, uint64_t* psuid, uint64_t* uid, int32_t* type) {
60,661✔
3483
  *psuid = pTableList->idInfo.suid;
60,661✔
3484
  *uid = pTableList->idInfo.uid;
60,661✔
3485
  *type = pTableList->idInfo.tableType;
60,661✔
3486
}
60,661✔
3487

3488
uint64_t tableListGetTableGroupId(const STableListInfo* pTableList, uint64_t tableUid) {
582,450,563✔
3489
  int32_t* slot = taosHashGet(pTableList->map, &tableUid, sizeof(tableUid));
582,450,563✔
3490
  if (slot == NULL) {
582,869,247✔
3491
    qDebug("table:%" PRIu64 " not found in table list", tableUid);
×
3492
    return -1;
×
3493
  }
3494

3495
  STableKeyInfo* pKeyInfo = taosArrayGet(pTableList->pTableList, *slot);
582,869,247✔
3496
  if (pKeyInfo == NULL) {
582,871,366✔
3497
    qDebug("table:%" PRIu64 " not found in table list", tableUid);
×
3498
    return -1;
×
3499
  }
3500
  return pKeyInfo->groupId;
582,871,366✔
3501
}
3502

3503
// TODO handle the group offset info, fix it, the rule of group output will be broken by this function
3504
// int32_t tableListRemoveTableInfo(STableListInfo* pTableList, uint64_t uid) {
3505
//   int32_t code = TSDB_CODE_SUCCESS;
3506
//   int32_t lino = 0;
3507

3508
//   int32_t* slot = taosHashGet(pTableList->map, &uid, sizeof(uid));
3509
//   if (slot == NULL) {
3510
//     qDebug("table:%" PRIu64 " not found in table list", uid);
3511
//     return 0;
3512
//   }
3513

3514
//   taosArrayRemove(pTableList->pTableList, *slot);
3515
//   code = taosHashRemove(pTableList->map, &uid, sizeof(uid));
3516

3517
//   _end:
3518
//   if (code != TSDB_CODE_SUCCESS) {
3519
//     qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
3520
//   } else {
3521
//     qDebug("uid:%" PRIu64 ", remove from table list", uid);
3522
//   }
3523

3524
//   return code;
3525
// }
3526

3527
int32_t tableListAddTableInfo(STableListInfo* pTableList, uint64_t uid, uint64_t gid) {
199,325✔
3528
  int32_t code = TSDB_CODE_SUCCESS;
199,325✔
3529
  int32_t lino = 0;
199,325✔
3530
  if (pTableList->map == NULL) {
199,325✔
3531
    pTableList->map = taosHashInit(32, taosGetDefaultHashFunction(TSDB_DATA_TYPE_BINARY), false, HASH_ENTRY_LOCK);
×
3532
    QUERY_CHECK_NULL(pTableList->map, code, lino, _end, terrno);
×
3533
  }
3534

3535
  STableKeyInfo keyInfo = {.uid = uid, .groupId = gid};
199,549✔
3536
  void*         p = taosHashGet(pTableList->map, &uid, sizeof(uid));
199,366✔
3537
  if (p != NULL) {
199,213✔
3538
    qInfo("table:%" PRId64 " already in tableIdList, ignore it", uid);
149✔
3539
    goto _end;
149✔
3540
  }
3541

3542
  void* tmp = taosArrayPush(pTableList->pTableList, &keyInfo);
199,064✔
3543
  QUERY_CHECK_NULL(tmp, code, lino, _end, terrno);
199,288✔
3544

3545
  int32_t slot = (int32_t)taosArrayGetSize(pTableList->pTableList) - 1;
199,288✔
3546
  code = taosHashPut(pTableList->map, &uid, sizeof(uid), &slot, sizeof(slot));
199,512✔
3547
  if (code != TSDB_CODE_SUCCESS) {
199,568✔
3548
    // we have checked the existence of uid in hash map above
3549
    QUERY_CHECK_CONDITION((code != TSDB_CODE_DUP_KEY), code, lino, _end, TSDB_CODE_QRY_EXECUTOR_INTERNAL_ERROR);
×
3550
    taosArrayPopTailBatch(pTableList->pTableList, 1);  // let's pop the last element in the array list
×
3551
  }
3552

3553
_end:
199,717✔
3554
  if (code != TSDB_CODE_SUCCESS) {
199,310✔
3555
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
3556
  } else {
3557
    qDebug("uid:%" PRIu64 ", groupId:%" PRIu64 " added into table list, slot:%d, total:%d", uid, gid, slot, slot + 1);
199,310✔
3558
  }
3559

3560
  return code;
199,646✔
3561
}
3562

3563
int32_t tableListGetGroupList(const STableListInfo* pTableList, int32_t ordinalGroupIndex, STableKeyInfo** pKeyInfo,
197,058,389✔
3564
                              int32_t* size) {
3565
  int32_t totalGroups = tableListGetOutputGroups(pTableList);
197,058,389✔
3566
  int32_t numOfTables = 0;
197,070,423✔
3567
  int32_t code = tableListGetSize(pTableList, &numOfTables);
197,121,316✔
3568
  if (code != TSDB_CODE_SUCCESS) {
197,048,111✔
3569
    qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(code));
×
3570
    return code;
×
3571
  }
3572

3573
  if (ordinalGroupIndex < 0 || ordinalGroupIndex >= totalGroups) {
197,048,111✔
3574
    return TSDB_CODE_INVALID_PARA;
×
3575
  }
3576

3577
  // here handle two special cases:
3578
  // 1. only one group exists, and 2. one table exists for each group.
3579
  if (totalGroups == 1) {
197,048,111✔
3580
    *size = numOfTables;
196,677,065✔
3581
    *pKeyInfo = (*size == 0) ? NULL : taosArrayGet(pTableList->pTableList, 0);
196,633,099✔
3582
    return TSDB_CODE_SUCCESS;
196,719,856✔
3583
  } else if (totalGroups == numOfTables) {
371,241✔
3584
    *size = 1;
330,415✔
3585
    *pKeyInfo = taosArrayGet(pTableList->pTableList, ordinalGroupIndex);
330,415✔
3586
    return TSDB_CODE_SUCCESS;
330,415✔
3587
  }
3588

3589
  int32_t offset = pTableList->groupOffset[ordinalGroupIndex];
41,172✔
3590
  if (ordinalGroupIndex < totalGroups - 1) {
55,260✔
3591
    *size = pTableList->groupOffset[ordinalGroupIndex + 1] - offset;
40,652✔
3592
  } else {
3593
    *size = numOfTables - offset;
14,608✔
3594
  }
3595

3596
  *pKeyInfo = taosArrayGet(pTableList->pTableList, offset);
55,260✔
3597
  return TSDB_CODE_SUCCESS;
55,260✔
3598
}
3599

3600
int32_t tableListGetOutputGroups(const STableListInfo* pTableList) { return pTableList->numOfOuputGroups; }
574,597,811✔
3601

3602
bool oneTableForEachGroup(const STableListInfo* pTableList) { return pTableList->oneTableForEachGroup; }
575,762✔
3603

3604
STableListInfo* tableListCreate() {
225,476,500✔
3605
  STableListInfo* pListInfo = taosMemoryCalloc(1, sizeof(STableListInfo));
225,476,500✔
3606
  if (pListInfo == NULL) {
225,201,159✔
3607
    return NULL;
×
3608
  }
3609

3610
  pListInfo->remainGroups = NULL;
225,201,159✔
3611
  pListInfo->pTableList = taosArrayInit(4, sizeof(STableKeyInfo));
225,240,530✔
3612
  if (pListInfo->pTableList == NULL) {
225,244,651✔
3613
    goto _error;
×
3614
  }
3615

3616
  pListInfo->map = taosHashInit(1024, taosGetDefaultHashFunction(TSDB_DATA_TYPE_BIGINT), false, HASH_ENTRY_LOCK);
225,368,137✔
3617
  if (pListInfo->map == NULL) {
225,634,024✔
3618
    goto _error;
×
3619
  }
3620

3621
  pListInfo->numOfOuputGroups = 1;
225,632,934✔
3622
  return pListInfo;
225,634,027✔
3623

3624
_error:
×
3625
  tableListDestroy(pListInfo);
×
3626
  return NULL;
×
3627
}
3628

3629
void tableListDestroy(STableListInfo* pTableListInfo) {
234,510,886✔
3630
  if (pTableListInfo == NULL) {
234,510,886✔
3631
    return;
9,100,209✔
3632
  }
3633

3634
  taosArrayDestroy(pTableListInfo->pTableList);
225,410,677✔
3635
  taosMemoryFreeClear(pTableListInfo->groupOffset);
225,190,655✔
3636

3637
  taosHashCleanup(pTableListInfo->map);
225,201,113✔
3638
  taosHashCleanup(pTableListInfo->remainGroups);
225,420,598✔
3639
  pTableListInfo->pTableList = NULL;
225,406,449✔
3640
  pTableListInfo->map = NULL;
225,474,408✔
3641
  taosMemoryFree(pTableListInfo);
225,456,619✔
3642
}
3643

3644
void tableListClear(STableListInfo* pTableListInfo) {
151,496✔
3645
  if (pTableListInfo == NULL) {
151,496✔
3646
    return;
×
3647
  }
3648

3649
  taosArrayClear(pTableListInfo->pTableList);
151,496✔
3650
  taosHashClear(pTableListInfo->map);
151,563✔
3651
  taosHashClear(pTableListInfo->remainGroups);
151,776✔
3652
  taosMemoryFree(pTableListInfo->groupOffset);
151,776✔
3653
  pTableListInfo->numOfOuputGroups = 1;
151,776✔
3654
  pTableListInfo->oneTableForEachGroup = false;
151,776✔
3655
}
3656

3657
static int32_t orderbyGroupIdComparFn(const void* p1, const void* p2) {
486,741,766✔
3658
  STableKeyInfo* pInfo1 = (STableKeyInfo*)p1;
486,741,766✔
3659
  STableKeyInfo* pInfo2 = (STableKeyInfo*)p2;
486,741,766✔
3660

3661
  if (pInfo1->groupId == pInfo2->groupId) {
486,741,766✔
3662
    return 0;
471,100,895✔
3663
  } else {
3664
    return pInfo1->groupId < pInfo2->groupId ? -1 : 1;
15,642,471✔
3665
  }
3666
}
3667

3668
int32_t sortTableGroup(STableListInfo* pTableListInfo) {
18,963,826✔
3669
  int32_t code = TSDB_CODE_SUCCESS;
18,963,826✔
3670
  taosArraySort(pTableListInfo->pTableList, orderbyGroupIdComparFn);
18,963,826✔
3671
  int32_t size = taosArrayGetSize(pTableListInfo->pTableList);
18,972,435✔
3672
  if (size == 0) {
18,967,792✔
3673
    pTableListInfo->numOfOuputGroups = 0;
×
3674
    return code;
×
3675
  }
3676

3677
  SArray* pList = taosArrayInit(4, sizeof(int32_t));
18,967,792✔
3678
  if (!pList) {
18,965,450✔
3679
    code = terrno;
×
3680
    goto end;
×
3681
  }
3682

3683
  STableKeyInfo* pInfo = taosArrayGet(pTableListInfo->pTableList, 0);
18,965,450✔
3684
  if (pInfo == NULL) {
18,963,245✔
3685
    qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
3686
    code = terrno;
×
3687
    goto end;
×
3688
  }
3689
  uint64_t gid = pInfo->groupId;
18,963,245✔
3690

3691
  int32_t start = 0;
18,966,920✔
3692
  void*   tmp = taosArrayPush(pList, &start);
18,966,140✔
3693
  if (tmp == NULL) {
18,966,140✔
3694
    qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
3695
    code = terrno;
×
3696
    goto end;
×
3697
  }
3698

3699
  for (int32_t i = 1; i < size; ++i) {
120,375,881✔
3700
    pInfo = taosArrayGet(pTableListInfo->pTableList, i);
101,412,170✔
3701
    if (pInfo == NULL) {
101,408,643✔
3702
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
3703
      code = terrno;
×
3704
      goto end;
×
3705
    }
3706
    if (pInfo->groupId != gid) {
101,408,643✔
3707
      tmp = taosArrayPush(pList, &i);
3,422,760✔
3708
      if (tmp == NULL) {
3,422,760✔
3709
        qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
3710
        code = terrno;
×
3711
        goto end;
×
3712
      }
3713
      gid = pInfo->groupId;
3,422,760✔
3714
    }
3715
  }
3716

3717
  pTableListInfo->numOfOuputGroups = taosArrayGetSize(pList);
18,971,459✔
3718
  pTableListInfo->groupOffset = taosMemoryMalloc(sizeof(int32_t) * pTableListInfo->numOfOuputGroups);
18,970,080✔
3719
  if (pTableListInfo->groupOffset == NULL) {
18,959,934✔
3720
    code = terrno;
×
3721
    goto end;
×
3722
  }
3723

3724
  memcpy(pTableListInfo->groupOffset, taosArrayGet(pList, 0), sizeof(int32_t) * pTableListInfo->numOfOuputGroups);
18,955,869✔
3725

3726
end:
18,963,634✔
3727
  taosArrayDestroy(pList);
18,965,407✔
3728
  return code;
18,949,906✔
3729
}
3730

3731
int32_t buildGroupIdMapForAllTables(STableListInfo* pTableListInfo, SReadHandle* pHandle, SScanPhysiNode* pScanNode,
207,994,154✔
3732
                                    SNodeList* group, bool groupSort, uint8_t* digest, SStorageAPI* pAPI, SHashObj* groupIdMap) {
3733
  int32_t code = TSDB_CODE_SUCCESS;
207,994,154✔
3734

3735
  bool   groupByTbname = groupbyTbname(group);
207,994,154✔
3736
  size_t numOfTables = taosArrayGetSize(pTableListInfo->pTableList);
207,980,201✔
3737
  if (!numOfTables) {
207,974,451✔
3738
    return code;
6,253✔
3739
  }
3740
  qDebug("numOfTables:%zu, groupByTbname:%d, group:%p", numOfTables, groupByTbname, group);
207,968,198✔
3741
  if (group == NULL || groupByTbname) {
207,906,606✔
3742
    if (tsCountAlwaysReturnValue && QUERY_NODE_PHYSICAL_PLAN_TABLE_SCAN == nodeType(pScanNode) &&
206,136,289✔
3743
        ((STableScanPhysiNode*)pScanNode)->needCountEmptyTable) {
178,599,749✔
3744
      pTableListInfo->remainGroups =
7,349,163✔
3745
          taosHashInit(numOfTables, taosGetDefaultHashFunction(TSDB_DATA_TYPE_BIGINT), false, HASH_NO_LOCK);
7,349,573✔
3746
      if (pTableListInfo->remainGroups == NULL) {
7,349,573✔
3747
        return terrno;
×
3748
      }
3749

3750
      for (int i = 0; i < numOfTables; i++) {
25,079,893✔
3751
        STableKeyInfo* info = taosArrayGet(pTableListInfo->pTableList, i);
17,730,454✔
3752
        if (!info) {
17,731,006✔
3753
          qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
3754
          return terrno;
×
3755
        }
3756
        info->groupId = groupByTbname ? info->uid : 0;
17,731,006✔
3757
        int32_t tempRes = taosHashPut(pTableListInfo->remainGroups, &(info->groupId), sizeof(info->groupId),
17,731,416✔
3758
                                      &(info->uid), sizeof(info->uid));
17,729,764✔
3759
        if (tempRes != TSDB_CODE_SUCCESS && tempRes != TSDB_CODE_DUP_KEY) {
17,730,868✔
3760
          qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(tempRes));
×
3761
          return tempRes;
×
3762
        }
3763
      }
3764
    } else {
3765
      for (int32_t i = 0; i < numOfTables; i++) {
672,609,939✔
3766
        STableKeyInfo* info = taosArrayGet(pTableListInfo->pTableList, i);
473,928,622✔
3767
        if (!info) {
473,915,619✔
3768
          qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
3769
          return terrno;
×
3770
        }
3771
        info->groupId = groupByTbname ? info->uid : 0;
473,915,619✔
3772
        
3773
      }
3774
    }
3775
    if (groupIdMap && group != NULL){
206,030,756✔
3776
      getColInfoResultForGroupbyForStream(pHandle->vnode, group, pTableListInfo, pAPI, groupIdMap);
106,748✔
3777
    }
3778

3779
    pTableListInfo->oneTableForEachGroup = groupByTbname;
206,030,756✔
3780
    if (numOfTables == 1 && pTableListInfo->idInfo.tableType == TSDB_CHILD_TABLE) {
206,108,211✔
3781
      pTableListInfo->oneTableForEachGroup = true;
92,101,898✔
3782
    }
3783

3784
    if (groupSort && groupByTbname) {
206,147,872✔
3785
      taosArraySort(pTableListInfo->pTableList, orderbyGroupIdComparFn);
1,441,136✔
3786
      pTableListInfo->numOfOuputGroups = numOfTables;
1,441,136✔
3787
    } else if (groupByTbname && pScanNode->groupOrderScan) {
204,706,736✔
3788
      pTableListInfo->numOfOuputGroups = numOfTables;
31,145✔
3789
    } else {
3790
      pTableListInfo->numOfOuputGroups = 1;
204,675,800✔
3791
    }
3792
    if (groupSort || pScanNode->groupOrderScan) {
206,253,575✔
3793
      code = sortTableGroup(pTableListInfo);
18,848,403✔
3794
    }
3795
  } else {
3796
    bool initRemainGroups = false;
1,770,317✔
3797
    if (QUERY_NODE_PHYSICAL_PLAN_TABLE_SCAN == nodeType(pScanNode)) {
1,770,317✔
3798
      STableScanPhysiNode* pTableScanNode = (STableScanPhysiNode*)pScanNode;
1,632,019✔
3799
      if (tsCountAlwaysReturnValue && pTableScanNode->needCountEmptyTable &&
1,632,019✔
3800
          !(groupSort || pScanNode->groupOrderScan)) {
832,951✔
3801
        initRemainGroups = true;
805,975✔
3802
      }
3803
    }
3804

3805
    code = getColInfoResultForGroupby(pHandle->vnode, group, pTableListInfo, digest, pAPI, initRemainGroups, groupIdMap);
1,770,317✔
3806
    if (code != TSDB_CODE_SUCCESS) {
1,770,797✔
3807
      return code;
×
3808
    }
3809

3810
    if (pScanNode->groupOrderScan) pTableListInfo->numOfOuputGroups = taosArrayGetSize(pTableListInfo->pTableList);
1,770,797✔
3811

3812
    if (groupSort || pScanNode->groupOrderScan) {
1,770,935✔
3813
      code = sortTableGroup(pTableListInfo);
152,453✔
3814
    }
3815
  }
3816

3817
  // add all table entry in the hash map
3818
  size_t size = taosArrayGetSize(pTableListInfo->pTableList);
207,956,266✔
3819
  for (int32_t i = 0; i < size; ++i) {
708,290,969✔
3820
    STableKeyInfo* p = taosArrayGet(pTableListInfo->pTableList, i);
500,248,384✔
3821
    if (!p) {
500,111,397✔
3822
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
3823
      return terrno;
×
3824
    }
3825
    int32_t tempRes = taosHashPut(pTableListInfo->map, &p->uid, sizeof(uint64_t), &i, sizeof(int32_t));
500,111,397✔
3826
    if (tempRes != TSDB_CODE_SUCCESS && tempRes != TSDB_CODE_DUP_KEY) {
500,327,228✔
3827
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(tempRes));
×
3828
      return tempRes;
×
3829
    }
3830
  }
3831

3832
  return code;
208,039,500✔
3833
}
3834

3835
int32_t createScanTableListInfo(SScanPhysiNode* pScanNode, SNodeList* pGroupTags, bool groupSort, SReadHandle* pHandle,
219,664,985✔
3836
                                STableListInfo* pTableListInfo, SNode* pTagCond, SNode* pTagIndexCond,
3837
                                SExecTaskInfo* pTaskInfo, SHashObj* groupIdMap) {
3838
  int64_t     st = taosGetTimestampUs();
219,679,281✔
3839
  const char* idStr = GET_TASKID(pTaskInfo);
219,679,281✔
3840

3841
  if (pHandle == NULL) {
219,409,650✔
3842
    qError("invalid handle, in creating operator tree, %s", idStr);
×
3843
    return TSDB_CODE_INVALID_PARA;
×
3844
  }
3845

3846
  if (pHandle->uid != 0) {
219,409,650✔
3847
    pScanNode->uid = pHandle->uid;
66,034✔
3848
    pScanNode->tableType = TSDB_CHILD_TABLE;
66,034✔
3849
  }
3850
  uint8_t digest[17] = {0};
219,614,938✔
3851
  int32_t code = getTableList(pHandle->vnode, pScanNode, pTagCond, pTagIndexCond, pTableListInfo, digest, idStr,
219,596,360✔
3852
                              &pTaskInfo->storageAPI, pTaskInfo->pStreamRuntimeInfo);
219,623,746✔
3853
  if (code != TSDB_CODE_SUCCESS) {
219,789,989✔
3854
    qError("failed to getTableList, code:%s", tstrerror(code));
1,102✔
3855
    return code;
1,102✔
3856
  }
3857

3858
  int32_t numOfTables = taosArrayGetSize(pTableListInfo->pTableList);
219,788,887✔
3859

3860
  int64_t st1 = taosGetTimestampUs();
219,799,197✔
3861
  pTaskInfo->cost.extractListTime = (st1 - st) / 1000.0;
219,799,197✔
3862
  qDebug("extract queried table list completed, %d tables, elapsed time:%.2f ms %s", numOfTables,
219,758,164✔
3863
         pTaskInfo->cost.extractListTime, idStr);
3864

3865
  if (numOfTables == 0) {
219,766,772✔
3866
    qDebug("no table qualified for query, %s", idStr);
11,847,405✔
3867
    return TSDB_CODE_SUCCESS;
11,847,405✔
3868
  }
3869

3870
  code = buildGroupIdMapForAllTables(pTableListInfo, pHandle, pScanNode, pGroupTags, groupSort, digest, &pTaskInfo->storageAPI, groupIdMap);
207,919,367✔
3871
  if (code != TSDB_CODE_SUCCESS) {
207,979,702✔
3872
    return code;
×
3873
  }
3874

3875
  pTaskInfo->cost.groupIdMapTime = (taosGetTimestampUs() - st1) / 1000.0;
208,019,858✔
3876
  qDebug("generate group id map completed, elapsed time:%.2f ms %s", pTaskInfo->cost.groupIdMapTime, idStr);
207,972,083✔
3877

3878
  return TSDB_CODE_SUCCESS;
207,923,508✔
3879
}
3880

3881
char* getStreamOpName(uint16_t opType) {
8,516,079✔
3882
  switch (opType) {
8,516,079✔
3883
    case QUERY_NODE_PHYSICAL_PLAN_STREAM_SCAN:
×
3884
      return "stream scan";
×
3885
    case QUERY_NODE_PHYSICAL_PLAN_PROJECT:
8,280,132✔
3886
      return "project";
8,280,132✔
3887
    case QUERY_NODE_PHYSICAL_PLAN_EXTERNAL_WINDOW:
235,947✔
3888
      return "external window";
235,947✔
3889
  }
3890
  return "error name";
×
3891
}
3892

3893
void printDataBlock(SSDataBlock* pBlock, const char* flag, const char* taskIdStr, int64_t qId) {
439,713,029✔
3894
  if (qDebugFlag & DEBUG_TRACE) {
439,713,029✔
3895
    if (!pBlock) {
38,446✔
3896
      qDebug("%" PRIx64 " %s %s %s: Block is Null", qId, taskIdStr, flag, __func__);
6,953✔
3897
      return;
6,953✔
3898
    } else if (pBlock->info.rows == 0) {
31,493✔
3899
      qDebug("%" PRIx64 " %s %s %s: Block is Empty. block type %d", qId, taskIdStr, flag, __func__, pBlock->info.type);
×
3900
      return;
×
3901
    }
3902
    
3903
    char*   pBuf = NULL;
31,493✔
3904
    int32_t code = dumpBlockData(pBlock, flag, &pBuf, taskIdStr, qId);
31,493✔
3905
    if (code == 0) {
31,493✔
3906
      qDebugL("%" PRIx64 " %s %s", qId, __func__, pBuf);
31,493✔
3907
      taosMemoryFree(pBuf);
31,493✔
3908
    }
3909
  }
3910
}
3911

3912
void printSpecDataBlock(SSDataBlock* pBlock, const char* flag, const char* opStr, const char* taskIdStr) {
×
3913
  if (!pBlock) {
×
3914
    qDebug("%s===stream===%s %s: Block is Null", taskIdStr, flag, opStr);
×
3915
    return;
×
3916
  } else if (pBlock->info.rows == 0) {
×
3917
    qDebug("%s===stream===%s %s: Block is Empty. block type %d.skey:%" PRId64 ",ekey:%" PRId64 ",version%" PRId64,
×
3918
           taskIdStr, flag, opStr, pBlock->info.type, pBlock->info.window.skey, pBlock->info.window.ekey,
3919
           pBlock->info.version);
3920
    return;
×
3921
  }
3922
  if (qDebugFlag & DEBUG_TRACE) {
×
3923
    char* pBuf = NULL;
×
3924
    char  flagBuf[64];
×
3925
    snprintf(flagBuf, sizeof(flagBuf), "%s %s", flag, opStr);
×
3926
    int32_t code = dumpBlockData(pBlock, flagBuf, &pBuf, taskIdStr, 0);
×
3927
    if (code == 0) {
×
3928
      qDebug("%s", pBuf);
×
3929
      taosMemoryFree(pBuf);
×
3930
    }
3931
  }
3932
}
3933

3934
TSKEY getStartTsKey(STimeWindow* win, const TSKEY* tsCols) { return tsCols == NULL ? win->skey : tsCols[0]; }
11,096,554✔
3935

3936
void updateTimeWindowInfo(SColumnInfoData* pColData, const STimeWindow* pWin, int64_t delta) {
2,147,483,647✔
3937
  int64_t* ts = (int64_t*)pColData->pData;
2,147,483,647✔
3938

3939
  int64_t duration = pWin->ekey > pWin->skey ? pWin->ekey - pWin->skey + delta : pWin->skey - pWin->ekey + delta;
2,147,483,647✔
3940
  ts[2] = duration;            // set the duration
2,147,483,647✔
3941
  ts[3] = pWin->skey;          // window start key
2,147,483,647✔
3942
  ts[4] = pWin->ekey + delta;  // window end key
2,147,483,647✔
3943
}
2,147,483,647✔
3944

3945
int32_t compKeys(const SArray* pSortGroupCols, const char* oldkeyBuf, int32_t oldKeysLen, const SSDataBlock* pBlock,
943,368,839✔
3946
                 int32_t rowIndex) {
3947
  SColumnDataAgg* pColAgg = NULL;
943,368,839✔
3948
  const char*     isNull = oldkeyBuf;
943,368,839✔
3949
  const char*     p = oldkeyBuf + sizeof(int8_t) * pSortGroupCols->size;
943,368,839✔
3950

3951
  for (int32_t i = 0; i < pSortGroupCols->size; ++i) {
2,147,483,647✔
3952
    const SColumn*         pCol = (SColumn*)TARRAY_GET_ELEM(pSortGroupCols, i);
1,462,808,926✔
3953
    const SColumnInfoData* pColInfoData = TARRAY_GET_ELEM(pBlock->pDataBlock, pCol->slotId);
1,463,361,352✔
3954
    if (pBlock->pBlockAgg) pColAgg = &pBlock->pBlockAgg[pCol->slotId];
1,463,424,650✔
3955

3956
    if (colDataIsNull(pColInfoData, pBlock->info.rows, rowIndex, pColAgg)) {
2,147,483,647✔
3957
      if (isNull[i] != 1) return 1;
101,097,384✔
3958
    } else {
3959
      if (isNull[i] != 0) return 1;
1,362,469,583✔
3960
      const char* val = colDataGetData(pColInfoData, rowIndex);
1,361,769,922✔
3961
      if (pCol->type == TSDB_DATA_TYPE_JSON) {
1,361,891,243✔
3962
        int32_t len = getJsonValueLen(val);
×
3963
        if (memcmp(p, val, len) != 0) return 1;
×
3964
        p += len;
×
3965
      } else if (IS_VAR_DATA_TYPE(pCol->type)) {
1,361,759,711✔
3966
        if (IS_STR_DATA_BLOB(pCol->type)) {
461,263,919✔
3967
          if (memcmp(p, val, blobDataTLen(val)) != 0) return 1;
×
3968
          p += blobDataTLen(val);
×
3969
        } else {
3970
          if (memcmp(p, val, varDataTLen(val)) != 0) return 1;
461,966,650✔
3971
          p += varDataTLen(val);
454,858,303✔
3972
        }
3973
      } else {
3974
        if (0 != memcmp(p, val, pCol->bytes)) return 1;
899,921,695✔
3975
        p += pCol->bytes;
882,489,608✔
3976
      }
3977
    }
3978
  }
3979
  if ((int32_t)(p - oldkeyBuf) != oldKeysLen) return 1;
917,701,461✔
3980
  return 0;
917,688,585✔
3981
}
3982

3983
int32_t buildKeys(char* keyBuf, const SArray* pSortGroupCols, const SSDataBlock* pBlock, int32_t rowIndex) {
25,201,367✔
3984
  uint32_t        colNum = pSortGroupCols->size;
25,201,367✔
3985
  SColumnDataAgg* pColAgg = NULL;
25,201,994✔
3986
  char*           isNull = keyBuf;
25,201,994✔
3987
  char*           p = keyBuf + sizeof(int8_t) * colNum;
25,201,994✔
3988

3989
  for (int32_t i = 0; i < colNum; ++i) {
75,295,159✔
3990
    const SColumn*         pCol = (SColumn*)TARRAY_GET_ELEM(pSortGroupCols, i);
50,091,075✔
3991
    const SColumnInfoData* pColInfoData = TARRAY_GET_ELEM(pBlock->pDataBlock, pCol->slotId);
50,094,001✔
3992
    if (pCol->slotId > pBlock->pDataBlock->size) continue;
50,093,448✔
3993

3994
    if (pBlock->pBlockAgg) pColAgg = &pBlock->pBlockAgg[pCol->slotId];
50,094,075✔
3995

3996
    if (colDataIsNull(pColInfoData, pBlock->info.rows, rowIndex, pColAgg)) {
100,190,584✔
3997
      isNull[i] = 1;
1,801,338✔
3998
    } else {
3999
      isNull[i] = 0;
48,294,544✔
4000
      const char* val = colDataGetData(pColInfoData, rowIndex);
48,295,589✔
4001
      if (pCol->type == TSDB_DATA_TYPE_JSON) {
48,296,634✔
4002
        int32_t len = getJsonValueLen(val);
×
4003
        memcpy(p, val, len);
×
4004
        p += len;
×
4005
      } else if (IS_VAR_DATA_TYPE(pCol->type)) {
48,291,409✔
4006
        if (IS_STR_DATA_BLOB(pCol->type)) {
7,262,483✔
4007
          blobDataCopy(p, val);
×
4008
          p += blobDataTLen(val);
×
4009
        } else {
4010
          varDataCopy(p, val);
7,263,737✔
4011
          p += varDataTLen(val);
7,263,737✔
4012
        }
4013
      } else {
4014
        memcpy(p, val, pCol->bytes);
41,031,643✔
4015
        p += pCol->bytes;
41,031,016✔
4016
      }
4017
    }
4018
  }
4019
  return (int32_t)(p - keyBuf);
25,204,084✔
4020
}
4021

4022
uint64_t calcGroupId(char* pData, int32_t len) {
2,147,483,647✔
4023
  T_MD5_CTX context;
2,147,483,647✔
4024
  tMD5Init(&context);
2,147,483,647✔
4025
  tMD5Update(&context, (uint8_t*)pData, len);
2,147,483,647✔
4026
  tMD5Final(&context);
2,147,483,647✔
4027

4028
  // NOTE: only extract the initial 8 bytes of the final MD5 digest
4029
  uint64_t id = 0;
2,147,483,647✔
4030
  memcpy(&id, context.digest, sizeof(uint64_t));
2,147,483,647✔
4031
  if (0 == id) memcpy(&id, context.digest + 8, sizeof(uint64_t));
2,147,483,647✔
4032
  return id;
2,147,483,647✔
4033
}
4034

4035
SNodeList* makeColsNodeArrFromSortKeys(SNodeList* pSortKeys) {
41,628✔
4036
  SNode*     node;
4037
  SNodeList* ret = NULL;
41,628✔
4038
  FOREACH(node, pSortKeys) {
126,800✔
4039
    SOrderByExprNode* pSortKey = (SOrderByExprNode*)node;
85,172✔
4040
    int32_t           code = nodesListMakeAppend(&ret, pSortKey->pExpr);
85,172✔
4041
    if (code != TSDB_CODE_SUCCESS) {
85,172✔
4042
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(code));
×
4043
      terrno = code;
×
4044
      return NULL;
×
4045
    }
4046
  }
4047
  return ret;
41,628✔
4048
}
4049

4050
int32_t extractKeysLen(const SArray* keys, int32_t* pLen) {
41,628✔
4051
  int32_t code = TSDB_CODE_SUCCESS;
41,628✔
4052
  int32_t lino = 0;
41,628✔
4053
  int32_t len = 0;
41,628✔
4054
  int32_t keyNum = taosArrayGetSize(keys);
41,628✔
4055
  for (int32_t i = 0; i < keyNum; ++i) {
106,160✔
4056
    SColumn* pCol = (SColumn*)taosArrayGet(keys, i);
64,532✔
4057
    QUERY_CHECK_NULL(pCol, code, lino, _end, terrno);
64,532✔
4058
    len += pCol->bytes;
64,532✔
4059
  }
4060
  len += sizeof(int8_t) * keyNum;  // null flag
41,628✔
4061
  *pLen = len;
41,628✔
4062

4063
_end:
41,628✔
4064
  if (code != TSDB_CODE_SUCCESS) {
41,628✔
4065
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
4066
  }
4067
  return code;
41,628✔
4068
}
4069

4070
int32_t parseErrorMsgFromAnalyticServer(SJson* pJson, const char* pId) {
×
4071
  int32_t code = TSDB_CODE_ANA_ANODE_RETURN_ERROR;
×
4072
  if (pJson == NULL) {
×
4073
    return code;
×
4074
  }
4075

4076
  char    pMsg[1024] = {0};
×
4077
  int32_t ret = tjsonGetStringValue(pJson, "msg", pMsg);
×
4078

4079
  if (ret == 0) {
×
4080
    qError("%s failed to exec imputation, msg:%s", pId, pMsg);
×
4081
    if (strstr(pMsg, "white noise") != NULL) {
×
4082
      code = TSDB_CODE_ANA_WN_DATA;
×
4083
    } else if (strstr(pMsg, "white-noise") != NULL) {
×
4084
      code = TSDB_CODE_ANA_WN_DATA;
×
4085
    } else if (strstr(pMsg, "[Errno 111] Connection refused") != NULL) {
×
4086
      code = TSDB_CODE_ANA_ALGO_NOT_LOAD;
×
4087
    }
4088
  } else {
4089
    qError("%s failed to extract msg from server, unknown error", pId);
×
4090
  }
4091

4092
  return code;
×
4093
}
4094

4095

4096
int32_t createBlockFromRemoteValueNode(SSDataBlock** ppBlock, SRemoteValueNode* pRemote) {
25,969,931✔
4097
  SValueNode* pVal = (SValueNode*)pRemote;
25,969,931✔
4098
  int32_t code = 0;
25,969,931✔
4099
  SSDataBlock* pBlock = taosMemoryCalloc(1, sizeof(SSDataBlock));
25,969,931✔
4100
  if (pBlock == NULL) {
25,967,202✔
4101
    return terrno;
×
4102
  }
4103

4104
  pBlock->pDataBlock = taosArrayInit(1, sizeof(SColumnInfoData));
25,967,202✔
4105
  if (pBlock->pDataBlock == NULL) {
25,966,024✔
4106
    code = terrno;
×
4107
    taosMemoryFree(pBlock);
×
4108
    return code;
×
4109
  }
4110

4111
  SColumnInfoData idata =
25,967,182✔
4112
      createColumnInfoData(pVal->node.resType.type, pVal->node.resType.bytes, 0);
25,969,371✔
4113
  idata.info.scale = pVal->node.resType.scale;
25,970,474✔
4114
  idata.info.precision = pVal->node.resType.precision;
25,969,872✔
4115

4116
  code = blockDataAppendColInfo(pBlock, &idata);
25,968,229✔
4117
  if (code != TSDB_CODE_SUCCESS) {
25,967,625✔
4118
    qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(code));
×
4119
    blockDataDestroy(pBlock);
×
4120
    *ppBlock = NULL;
×
4121
    return code;
×
4122
  }
4123

4124
  *ppBlock = pBlock;
25,967,625✔
4125

4126
  return code;
25,969,329✔
4127
}
4128

4129

4130
int32_t extractSingleRspBlock(SRetrieveTableRsp* pRetrieveRsp, SSDataBlock* pb) {
25,969,872✔
4131
  int32_t            code = TSDB_CODE_SUCCESS;
25,969,872✔
4132
  int32_t            lino = 0;
25,969,872✔
4133
  void*              decompBuf = NULL;
25,969,872✔
4134

4135
  char* pNextStart = pRetrieveRsp->data;
25,969,872✔
4136
  char* pStart = pNextStart;
25,968,246✔
4137

4138
  int32_t index = 0;
25,968,730✔
4139

4140
  if (pRetrieveRsp->compressed) {  // decompress the data
25,968,730✔
4141
    decompBuf = taosMemoryMalloc(pRetrieveRsp->payloadLen);
×
4142
    QUERY_CHECK_NULL(decompBuf, code, lino, _end, terrno);
×
4143
  }
4144

4145
  int32_t compLen = *(int32_t*)pStart;
25,969,914✔
4146
  pStart += sizeof(int32_t);
25,969,872✔
4147

4148
  int32_t rawLen = *(int32_t*)pStart;
25,969,872✔
4149
  pStart += sizeof(int32_t);
25,969,872✔
4150
  QUERY_CHECK_CONDITION((compLen <= rawLen && compLen != 0), code, lino, _end, TSDB_CODE_QRY_EXECUTOR_INTERNAL_ERROR);
25,970,474✔
4151

4152
  pNextStart = pStart + compLen;
25,970,474✔
4153
  if (pRetrieveRsp->compressed && (compLen < rawLen)) {
25,964,808✔
4154
    int32_t t = tsDecompressString(pStart, compLen, 1, decompBuf, rawLen, ONE_STAGE_COMP, NULL, 0);
×
4155
    QUERY_CHECK_CONDITION((t == rawLen), code, lino, _end, TSDB_CODE_QRY_EXECUTOR_INTERNAL_ERROR);
×
4156
    pStart = decompBuf;
×
4157
  }
4158

4159
  code = blockDecodeInternal(pb, pStart, (const char**)&pStart);
25,969,872✔
4160
  if (code != 0) {
25,963,671✔
4161
    taosMemoryFreeClear(pRetrieveRsp);
×
4162
    goto _end;
×
4163
  }
4164

4165
_end:
25,963,671✔
4166
  if (code != TSDB_CODE_SUCCESS) {
25,967,624✔
4167
    blockDataDestroy(pb);
×
4168
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
4169
  }
4170
  return code;
25,967,624✔
4171
}
4172

4173
int32_t setValueFromResBlock(STaskSubJobCtx* ctx, SRemoteValueNode* pRes, SSDataBlock* pBlock) {
25,968,167✔
4174
  int32_t code = 0;
25,968,167✔
4175
  bool needFree = true;
25,968,167✔
4176
  int32_t colNum = taosArrayGetSize(pBlock->pDataBlock);
25,967,607✔
4177
  if (NULL == pBlock->pDataBlock || 1 != colNum || pBlock->info.rows > 1) {
25,966,521✔
4178
    qError("%s invalid scl fetch res block, pDataBlock:%p, colNum:%d, rows:%" PRId64, 
4,498✔
4179
      ctx->idStr, pBlock->pDataBlock, colNum, pBlock->info.rows);
4180
    return TSDB_CODE_PAR_INVALID_SCALAR_SUBQ_RES_ROWS;
×
4181
  }
4182
  
4183
  pRes->val.node.type = QUERY_NODE_VALUE;
25,962,622✔
4184
  pRes->val.flag &= (~VALUE_FLAG_VAL_UNSET);
25,968,167✔
4185
  pRes->val.translate = true;
25,967,199✔
4186
  
4187
  SColumnInfoData* pCol = taosArrayGet(pBlock->pDataBlock, 0);
25,967,005✔
4188
  if (colDataIsNull_s(pCol, 0)) {
25,968,710✔
4189
    pRes->val.isNull = true;
2,127,674✔
4190
  } else {
4191
    code = nodesSetValueNodeValueExt(&pRes->val, colDataGetData(pCol, 0), &needFree);
23,841,036✔
4192
  }
4193

4194
  if (!needFree) {
25,964,231✔
4195
    pCol->pData = NULL;
17,304✔
4196
  }
4197

4198
  return code;
25,964,231✔
4199
}
4200

4201
int32_t remoteFetchCallBack(void* param, SDataBuf* pMsg, int32_t code) {
46,899,340✔
4202
  SScalarFetchParam* pParam = (SScalarFetchParam*)param;
46,899,340✔
4203
  STaskSubJobCtx* ctx = pParam->pSubJobCtx;
46,899,340✔
4204
  SSDataBlock* pResBlock = NULL;
46,898,800✔
4205
  
4206
  taosMemoryFreeClear(pMsg->pEpSet);
46,895,975✔
4207

4208
  if (NULL == ctx) {
46,895,894✔
4209
    qWarn("scl fetch ctx not exists since it may have been released");
6,020✔
4210
    goto _exit;
6,020✔
4211
  }
4212

4213
  qDebug("%s subQIdx %d got rsp, code:%d, rsp:%p", ctx->idStr, pParam->subQIdx, code, pMsg->pData);
46,889,874✔
4214

4215
  taosWLockLatch(&ctx->lock);
46,889,874✔
4216
  ctx->param = NULL;
46,888,965✔
4217
  taosWUnLockLatch(&ctx->lock);
46,888,965✔
4218

4219
  if (ctx->transporterId > 0) {
46,892,777✔
4220
    int32_t ret = asyncFreeConnById(ctx->rpcHandle, ctx->transporterId);
46,893,320✔
4221
    if (ret != 0) {
46,893,320✔
4222
      qDebug("%s failed to free subQ rpc handle, code:%s, subQIdx:%d", ctx->idStr, tstrerror(ret), pParam->subQIdx);
×
4223
    }
4224
    ctx->transporterId = -1;
46,893,320✔
4225
  }
4226

4227
  if (0 == code && NULL == pMsg->pData) {
46,892,777✔
4228
    qError("%s invalid rsp msg, msgType:%d, len:%d", ctx->idStr, pMsg->msgType, pMsg->len);
×
4229
    code = TSDB_CODE_QRY_INVALID_MSG;
×
4230
  }
4231

4232
  if (code == TSDB_CODE_SUCCESS) {
46,893,320✔
4233
    SRetrieveTableRsp* pRsp = pMsg->pData;
38,676,050✔
4234
    pRsp->numOfRows = htobe64(pRsp->numOfRows);
38,676,050✔
4235
    pRsp->compLen = htonl(pRsp->compLen);
38,676,050✔
4236
    pRsp->payloadLen = htonl(pRsp->payloadLen);
38,673,688✔
4237
    pRsp->numOfCols = htonl(pRsp->numOfCols);
38,669,859✔
4238
    pRsp->useconds = htobe64(pRsp->useconds);
38,667,329✔
4239
    pRsp->numOfBlocks = htonl(pRsp->numOfBlocks);
38,667,270✔
4240

4241
    if (pRsp->numOfRows > 1 || pRsp->numOfBlocks > 1 || !pRsp->completed) {
38,671,000✔
4242
      qError("%s invalid scl fetch rsp received, subQIdx:%d, rows:%" PRId64 ", blocks:%d, completed:%d", 
898,674✔
4243
        ctx->idStr, pParam->subQIdx, pRsp->numOfRows, pRsp->numOfBlocks, pRsp->completed);
4244
      ctx->code = TSDB_CODE_PAR_INVALID_SCALAR_SUBQ_RES_ROWS;
898,674✔
4245
    } else if (0 == pRsp->numOfRows) {
37,770,350✔
4246
      SRemoteValueNode* pRemote = (SRemoteValueNode*)pParam->pRes;
11,806,073✔
4247
      pRemote->val.node.type = QUERY_NODE_VALUE;
11,806,613✔
4248
      pRemote->val.isNull = true;
11,807,693✔
4249
      pRemote->val.translate = true;
11,807,153✔
4250
      pRemote->val.flag &= (~VALUE_FLAG_VAL_UNSET);
11,807,693✔
4251
      taosArraySet(ctx->subResValues, pParam->subQIdx, &pParam->pRes);
11,807,693✔
4252
    } else {
4253
      qDebug("%s scl fetch rsp received, subQIdx:%d, rows:%" PRId64 , ctx->idStr, pParam->subQIdx, pRsp->numOfRows);
25,966,446✔
4254
      ctx->code = createBlockFromRemoteValueNode(&pResBlock, pParam->pRes);
25,966,446✔
4255
      if (TSDB_CODE_SUCCESS == ctx->code) {
25,970,474✔
4256
        ctx->code = blockDataEnsureCapacity(pResBlock, 1);
25,966,480✔
4257
      }
4258
      if (TSDB_CODE_SUCCESS == ctx->code) {
25,968,828✔
4259
        ctx->code = extractSingleRspBlock(pRsp, pResBlock);
25,968,769✔
4260
      }
4261
      if (TSDB_CODE_SUCCESS == ctx->code) {
25,965,995✔
4262
        ctx->code = setValueFromResBlock(ctx, pParam->pRes, pResBlock);
25,967,666✔
4263
      }
4264
      if (TSDB_CODE_SUCCESS == ctx->code) {
25,965,855✔
4265
        taosArraySet(ctx->subResValues, pParam->subQIdx, &pParam->pRes);
25,969,329✔
4266
      }
4267
    }
4268
  } else {
4269
    ctx->code = rpcCvtErrCode(code);
8,217,270✔
4270
    if (ctx->code != code) {
8,217,270✔
4271
      qError("%s scl fetch rsp received, subQIdx:%d, error:%s, cvted error: %s", ctx->idStr, pParam->subQIdx,
×
4272
             tstrerror(code), tstrerror(ctx->code));
4273
    } else {
4274
      qError("%s scl fetch rsp received, subQIdx:%d, error:%s", ctx->idStr, pParam->subQIdx, tstrerror(code));
8,217,270✔
4275
    }
4276
  }
4277
  
4278
  code = tsem_post(&pParam->pSubJobCtx->ready);
46,889,446✔
4279
  if (code != TSDB_CODE_SUCCESS) {
46,892,777✔
4280
    qError("failed to invoke post when scl fetch rsp is ready, code:%s", tstrerror(code));
×
4281
  }
4282

4283
_exit:
46,898,797✔
4284

4285
  taosMemoryFree(pMsg->pData);
46,898,198✔
4286
  blockDataDestroy(pResBlock);
46,896,575✔
4287

4288
  return code;
46,898,257✔
4289
}
4290

4291

4292
int32_t fetchRemoteValueImpl(STaskSubJobCtx* ctx, int32_t subQIdx, SRemoteValueNode* pRes) {
46,881,036✔
4293
  int32_t          code = TSDB_CODE_SUCCESS;
46,881,036✔
4294
  int32_t          lino = 0;
46,881,036✔
4295
  SDownstreamSourceNode* pSource = (SDownstreamSourceNode*)taosArrayGetP(ctx->subEndPoints, subQIdx);
46,881,036✔
4296

4297
  SResFetchReq req = {0};
46,879,601✔
4298
  req.header.vgId = pSource->addr.nodeId;
46,874,124✔
4299
  req.sId = pSource->sId;
46,882,440✔
4300
  req.clientId = pSource->clientId;
46,888,617✔
4301
  req.taskId = pSource->taskId;
46,868,946✔
4302
  req.queryId = ctx->queryId;
46,874,728✔
4303
  req.execId = pSource->execId;
46,873,592✔
4304

4305
  int32_t msgSize = tSerializeSResFetchReq(NULL, 0, &req, false);
46,832,063✔
4306
  if (msgSize < 0) {
46,862,623✔
4307
    return msgSize;
×
4308
  }
4309

4310
  void* msg = taosMemoryCalloc(1, msgSize);
46,862,623✔
4311
  if (NULL == msg) {
46,833,347✔
4312
    return terrno;
×
4313
  }
4314

4315
  msgSize = tSerializeSResFetchReq(msg, msgSize, &req, false);
46,833,347✔
4316
  if (msgSize < 0) {
46,866,251✔
4317
    taosMemoryFree(msg);
×
4318
    return msgSize;
×
4319
  }
4320

4321
  qDebug("%s scl build fetch msg and send to nodeId:%d, ep:%s, clientId:0x%" PRIx64 " taskId:0x%" PRIx64
46,866,251✔
4322
         ", execId:%d",
4323
         ctx->idStr, pSource->addr.nodeId, pSource->addr.epSet.eps[0].fqdn, pSource->clientId,
4324
         pSource->taskId, pSource->execId);
4325

4326
  // send the fetch remote task result reques
4327
  SMsgSendInfo* pMsgSendInfo = taosMemoryCalloc(1, sizeof(SMsgSendInfo));
46,878,761✔
4328
  if (NULL == pMsgSendInfo) {
46,871,999✔
4329
    taosMemoryFreeClear(msg);
×
4330
    qError("%s prepare message %d failed", ctx->idStr, (int32_t)sizeof(SMsgSendInfo));
×
4331
    return terrno;
×
4332
  }
4333

4334
  SScalarFetchParam* param = taosMemoryMalloc(sizeof(SScalarFetchParam));
46,871,999✔
4335
  if (NULL == param) {
46,848,692✔
4336
    taosMemoryFreeClear(msg);
×
4337
    taosMemoryFreeClear(pMsgSendInfo);
×
4338
    qError("%s prepare param %d failed", ctx->idStr, (int32_t)sizeof(SScalarFetchParam));
×
4339
    return terrno;
×
4340
  }
4341

4342
  taosWLockLatch(&ctx->lock);
46,848,692✔
4343
  
4344
  if (ctx->code) {
46,886,607✔
4345
    qError("task has been killed, error:%s", tstrerror(ctx->code));
×
4346
    taosMemoryFree(param);
×
4347
    code = ctx->code;
×
4348
    goto _end;
×
4349
  } else {
4350
    ctx->param = param;
46,869,571✔
4351
  }
4352
  
4353
  taosWUnLockLatch(&ctx->lock);
46,880,092✔
4354

4355
  param->subQIdx = subQIdx;
46,894,415✔
4356
  param->pRes = pRes;
46,894,415✔
4357
  param->pSubJobCtx = ctx;
46,895,468✔
4358

4359
  pMsgSendInfo->param = param;
46,878,806✔
4360
  pMsgSendInfo->paramFreeFp = taosAutoMemoryFree;
46,877,056✔
4361
  pMsgSendInfo->msgInfo.pData = msg;
46,889,782✔
4362
  pMsgSendInfo->msgInfo.len = msgSize;
46,869,854✔
4363
  pMsgSendInfo->msgType = pSource->fetchMsgType;
46,853,997✔
4364
  pMsgSendInfo->fp = remoteFetchCallBack;
46,874,216✔
4365
  pMsgSendInfo->requestId = ctx->queryId;
46,852,949✔
4366

4367
  code = asyncSendMsgToServer(ctx->rpcHandle, &pSource->addr.epSet, &ctx->transporterId, pMsgSendInfo);
46,861,380✔
4368
  QUERY_CHECK_CODE(code, lino, _end);
46,894,452✔
4369

4370
  code = qSemWait(ctx->pTaskInfo, &ctx->ready);
46,894,452✔
4371
  if (isTaskKilled(ctx->pTaskInfo)) {
46,899,444✔
4372
    code = getTaskCode(ctx->pTaskInfo);
11,438✔
4373
  } else {
4374
    code = ctx->code;
46,889,106✔
4375
  }
4376
      
4377
_end:
46,899,984✔
4378

4379
  taosWLockLatch(&ctx->lock);
46,899,984✔
4380
  ctx->param = NULL;
46,899,444✔
4381
  taosWUnLockLatch(&ctx->lock);
46,899,984✔
4382

4383
  if (code != TSDB_CODE_SUCCESS) {
46,900,544✔
4384
    qError("%s %s failed at line %d since %s", ctx->idStr, __func__, lino, tstrerror(code));
9,121,837✔
4385
  }
4386
  return code;
46,899,444✔
4387
}
4388

4389

4390
int32_t qFetchRemoteValue(void* pCtx, int32_t subQIdx, SRemoteValueNode* pRes) {
48,434,037✔
4391
  STaskSubJobCtx*  ctx = (STaskSubJobCtx*)pCtx;
48,434,037✔
4392
  int32_t code = 0, lino = 0;
48,434,037✔
4393
  int32_t       subEndPoinsNum = taosArrayGetSize(ctx->subEndPoints);
48,434,037✔
4394
  if (subQIdx >= subEndPoinsNum) {
48,410,679✔
4395
    qError("%s invalid subQIdx %d, subEndPointsNum:%d", ctx->idStr, subQIdx, subEndPoinsNum);
×
4396
    return TSDB_CODE_QRY_SUBQ_NOT_FOUND;
×
4397
  }
4398

4399
  SValueNode** ppRes = taosArrayGet(ctx->subResValues, subQIdx);
48,410,679✔
4400
  if (NULL == *ppRes) {
48,436,071✔
4401
    TAOS_CHECK_EXIT(fetchRemoteValueImpl(ctx, subQIdx, pRes));
46,891,929✔
4402
    *ppRes = (SValueNode*)pRes;
37,777,067✔
4403
  } else {
4404
    TAOS_CHECK_EXIT(valueNodeCopy(*ppRes, &pRes->val));
1,554,087✔
4405
    pRes->val.node.type = QUERY_NODE_VALUE;
1,554,087✔
4406
  }
4407

4408
_exit:
48,452,991✔
4409

4410
  if (code) {
48,452,991✔
4411
    qError("%s %s failed at line %d since %s", ctx->idStr, __func__, lino, tstrerror(code));
9,121,837✔
4412
  }
4413

4414
  return code;
48,453,531✔
4415
}
4416

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