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

taosdata / TDengine / #4894

22 Dec 2025 09:33AM UTC coverage: 65.72% (+0.2%) from 65.57%
#4894

push

travis-ci

web-flow
Update README.md (#34007)

* Update README.md

* docs: update table of contents and improve installation instructions in README

* docs: adjust words

---------

Signed-off-by: WANG Xu <feici02@outlook.com>
Co-authored-by: WANG Xu <feici02@outlook.com>

184394 of 280577 relevant lines covered (65.72%)

111859687.16 hits per line

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

75.48
/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

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

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

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

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

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

64
static int64_t getLimit(const SNode* pLimit) {
926,823,290✔
65
  return (NULL == pLimit || NULL == ((SLimitNode*)pLimit)->limit) ? -1 : ((SLimitNode*)pLimit)->limit->datum.i;
926,823,290✔
66
}
67
static int64_t getOffset(const SNode* pLimit) {
926,779,871✔
68
  return (NULL == pLimit || NULL == ((SLimitNode*)pLimit)->offset) ? -1 : ((SLimitNode*)pLimit)->offset->datum.i;
926,779,871✔
69
}
70
static void releaseColInfoData(void* pCol);
71

72
void initResultRowInfo(SResultRowInfo* pResultRowInfo) {
393,730,216✔
73
  pResultRowInfo->size = 0;
393,730,216✔
74
  pResultRowInfo->cur.pageId = -1;
393,759,451✔
75
}
393,817,452✔
76

77
void closeResultRow(SResultRow* pResultRow) { pResultRow->closed = true; }
6,096,877✔
78

79
void resetResultRow(SResultRow* pResultRow, size_t entrySize) {
404,011,924✔
80
  pResultRow->numOfRows = 0;
404,011,924✔
81
  pResultRow->closed = false;
404,015,515✔
82
  pResultRow->endInterp = false;
404,016,460✔
83
  pResultRow->startInterp = false;
404,018,161✔
84

85
  if (entrySize > 0) {
404,018,539✔
86
    memset(pResultRow->pEntryInfo, 0, entrySize);
404,017,783✔
87
  }
88
}
404,025,532✔
89

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

95
size_t getResultRowSize(SqlFunctionCtx* pCtx, int32_t numOfOutput) {
222,300,454✔
96
  int32_t rowSize = (numOfOutput * sizeof(SResultRowEntryInfo)) + sizeof(SResultRow);
222,300,454✔
97

98
  for (int32_t i = 0; i < numOfOutput; ++i) {
867,714,248✔
99
    rowSize += pCtx[i].resDataInfo.interBufSize;
645,432,019✔
100
  }
101

102
  return rowSize;
222,282,229✔
103
}
104

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

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

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

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

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

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

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

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

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

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

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

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

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

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

218
void cleanupGroupResInfo(SGroupResInfo* pGroupResInfo) {
102,686,698✔
219
  taosMemoryFreeClear(pGroupResInfo->pBuf);
102,686,698✔
220
  if (pGroupResInfo->freeItem) {
102,690,334✔
221
    //    taosArrayDestroy(pGroupResInfo->pRows);
222
    taosArrayDestroyEx(pGroupResInfo->pRows, freeEx);
×
223
    pGroupResInfo->freeItem = false;
×
224
    pGroupResInfo->pRows = NULL;
×
225
  } else {
226
    taosArrayDestroy(pGroupResInfo->pRows);
102,686,278✔
227
    pGroupResInfo->pRows = NULL;
102,687,074✔
228
  }
229
  pGroupResInfo->index = 0;
102,687,942✔
230
  pGroupResInfo->delIndex = 0;
102,687,942✔
231
}
102,687,142✔
232

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

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

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

251
static int32_t resultrowComparDesc(const void* p1, const void* p2) { return resultrowComparAsc(p2, p1); }
2,009,943,437✔
252

253
int32_t initGroupedResultInfo(SGroupResInfo* pGroupResInfo, SSHashObj* pHashmap, int32_t order) {
66,834,710✔
254
  int32_t code = TSDB_CODE_SUCCESS;
66,834,710✔
255
  int32_t lino = 0;
66,834,710✔
256
  if (pGroupResInfo->pRows != NULL) {
66,834,710✔
257
    taosArrayDestroy(pGroupResInfo->pRows);
3,280,915✔
258
  }
259
  if (pGroupResInfo->pBuf) {
66,842,183✔
260
    taosMemoryFree(pGroupResInfo->pBuf);
3,280,915✔
261
    pGroupResInfo->pBuf = NULL;
3,280,915✔
262
  }
263

264
  // extract the result rows information from the hash map
265
  int32_t size = tSimpleHashGetSize(pHashmap);
66,843,640✔
266

267
  void* pData = NULL;
66,831,385✔
268
  pGroupResInfo->pRows = taosArrayInit(size, POINTER_BYTES);
66,831,385✔
269
  QUERY_CHECK_NULL(pGroupResInfo->pRows, code, lino, _end, terrno);
66,846,849✔
270

271
  size_t  keyLen = 0;
66,838,317✔
272
  int32_t iter = 0;
66,837,579✔
273
  int64_t bufLen = 0, offset = 0;
66,840,608✔
274

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

281
  pGroupResInfo->pBuf = taosMemoryMalloc(bufLen);
66,807,042✔
282
  QUERY_CHECK_NULL(pGroupResInfo->pBuf, code, lino, _end, terrno);
66,848,155✔
283

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

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

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

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

299
  if (order == TSDB_ORDER_ASC || order == TSDB_ORDER_DESC) {
66,836,390✔
300
    __compar_fn_t fn = (order == TSDB_ORDER_ASC) ? resultrowComparAsc : resultrowComparDesc;
13,085,562✔
301
    size = POINTER_BYTES;
13,085,562✔
302
    taosSort(pGroupResInfo->pRows->pData, taosArrayGetSize(pGroupResInfo->pRows), size, fn);
13,085,562✔
303
  }
304

305
  pGroupResInfo->index = 0;
66,836,390✔
306

307
_end:
66,823,973✔
308
  if (code != TSDB_CODE_SUCCESS) {
66,845,938✔
309
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
310
  }
311
  return code;
66,845,938✔
312
}
313

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

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

325
bool hasRemainResults(SGroupResInfo* pGroupResInfo) {
232,774,370✔
326
  if (pGroupResInfo->pRows == NULL) {
232,774,370✔
327
    return false;
×
328
  }
329

330
  return pGroupResInfo->index < taosArrayGetSize(pGroupResInfo->pRows);
232,780,787✔
331
}
332

333
int32_t getNumOfTotalRes(SGroupResInfo* pGroupResInfo) {
125,628,887✔
334
  if (pGroupResInfo->pRows == 0) {
125,628,887✔
335
    return 0;
×
336
  }
337

338
  return (int32_t)taosArrayGetSize(pGroupResInfo->pRows);
125,633,508✔
339
}
340

341
SArray* createSortInfo(SNodeList* pNodeList) {
41,915,227✔
342
  size_t numOfCols = 0;
41,915,227✔
343

344
  if (pNodeList != NULL) {
41,915,227✔
345
    numOfCols = LIST_LENGTH(pNodeList);
41,866,096✔
346
  } else {
347
    numOfCols = 0;
49,431✔
348
  }
349

350
  SArray* pList = taosArrayInit(numOfCols, sizeof(SBlockOrderInfo));
41,915,925✔
351
  if (pList == NULL) {
41,918,701✔
352
    return pList;
×
353
  }
354

355
  for (int32_t i = 0; i < numOfCols; ++i) {
92,666,204✔
356
    SOrderByExprNode* pSortKey = (SOrderByExprNode*)nodesListGetNode(pNodeList, i);
50,740,366✔
357
    if (!pSortKey) {
50,745,024✔
358
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
359
      taosArrayDestroy(pList);
×
360
      pList = NULL;
×
361
      terrno = TSDB_CODE_QRY_EXECUTOR_INTERNAL_ERROR;
×
362
      break;
×
363
    }
364
    SBlockOrderInfo bi = {0};
50,745,024✔
365
    bi.order = (pSortKey->order == ORDER_ASC) ? TSDB_ORDER_ASC : TSDB_ORDER_DESC;
50,738,287✔
366
    bi.nullFirst = (pSortKey->nullOrder == NULL_ORDER_FIRST);
50,744,219✔
367

368
    if (nodeType(pSortKey->pExpr) != QUERY_NODE_COLUMN) {
50,744,999✔
369
      qError("invalid order by expr type:%d", nodeType(pSortKey->pExpr));
×
370
      taosArrayDestroy(pList);
×
371
      pList = NULL;
×
372
      terrno = TSDB_CODE_QRY_EXECUTOR_INTERNAL_ERROR;
×
373
      break;
×
374
    }
375
    
376
    SColumnNode* pColNode = (SColumnNode*)pSortKey->pExpr;
50,733,717✔
377
    bi.slotId = pColNode->slotId;
50,738,174✔
378
    void* tmp = taosArrayPush(pList, &bi);
50,748,593✔
379
    if (!tmp) {
50,748,593✔
380
      taosArrayDestroy(pList);
×
381
      pList = NULL;
×
382
      break;
×
383
    }
384
  }
385

386
  return pList;
41,926,437✔
387
}
388

389
SSDataBlock* createDataBlockFromDescNode(void* p) {
570,438,858✔
390
  SDataBlockDescNode* pNode = (SDataBlockDescNode*)p;
570,438,858✔
391
  int32_t      numOfCols = LIST_LENGTH(pNode->pSlots);
570,438,858✔
392
  SSDataBlock* pBlock = NULL;
570,482,965✔
393
  int32_t      code = createDataBlock(&pBlock);
570,482,701✔
394
  if (code) {
570,363,918✔
395
    terrno = code;
×
396
    return NULL;
×
397
  }
398

399
  pBlock->info.id.blockId = pNode->dataBlockId;
570,363,918✔
400
  pBlock->info.type = STREAM_INVALID;
570,374,803✔
401
  pBlock->info.calWin = (STimeWindow){.skey = INT64_MIN, .ekey = INT64_MAX};
570,398,788✔
402
  pBlock->info.watermark = INT64_MIN;
570,461,019✔
403

404
  for (int32_t i = 0; i < numOfCols; ++i) {
2,147,483,647✔
405
    SSlotDescNode* pDescNode = (SSlotDescNode*)nodesListGetNode(pNode->pSlots, i);
2,012,635,388✔
406
    if (!pDescNode) {
2,012,178,420✔
407
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(code));
×
408
      blockDataDestroy(pBlock);
×
409
      pBlock = NULL;
×
410
      terrno = TSDB_CODE_INVALID_PARA;
×
411
      break;
×
412
    }
413
    SColumnInfoData idata =
2,012,122,444✔
414
        createColumnInfoData(pDescNode->dataType.type, pDescNode->dataType.bytes, pDescNode->slotId);
2,012,392,486✔
415
    idata.info.scale = pDescNode->dataType.scale;
2,012,774,008✔
416
    idata.info.precision = pDescNode->dataType.precision;
2,012,825,826✔
417
    idata.info.noData = pDescNode->reserve;
2,012,852,682✔
418

419
    code = blockDataAppendColInfo(pBlock, &idata);
2,012,870,699✔
420
    if (code != TSDB_CODE_SUCCESS) {
2,012,693,768✔
421
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(code));
64,679✔
422
      blockDataDestroy(pBlock);
64,679✔
423
      pBlock = NULL;
×
424
      terrno = code;
×
425
      break;
×
426
    }
427
  }
428

429
  return pBlock;
570,623,833✔
430
}
431

432
int32_t prepareDataBlockBuf(SSDataBlock* pDataBlock, SColMatchInfo* pMatchInfo) {
205,877,266✔
433
  SDataBlockInfo* pBlockInfo = &pDataBlock->info;
205,877,266✔
434

435
  for (int32_t i = 0; i < taosArrayGetSize(pMatchInfo->pList); ++i) {
956,401,274✔
436
    SColMatchItem* pItem = taosArrayGet(pMatchInfo->pList, i);
757,520,815✔
437
    if (!pItem) {
757,503,378✔
438
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
439
      return terrno;
×
440
    }
441

442
    if (pItem->isPk) {
757,503,378✔
443
      SColumnInfoData* pInfoData = taosArrayGet(pDataBlock->pDataBlock, pItem->dstSlotId);
7,060,055✔
444
      if (!pInfoData) {
7,113,568✔
445
        qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
446
        return terrno;
×
447
      }
448
      pBlockInfo->pks[0].type = pInfoData->info.type;
7,113,568✔
449
      pBlockInfo->pks[1].type = pInfoData->info.type;
7,114,954✔
450

451
      // allocate enough buffer size, which is pInfoData->info.bytes
452
      if (IS_VAR_DATA_TYPE(pItem->dataType.type)) {
7,116,168✔
453
        pBlockInfo->pks[0].pData = taosMemoryCalloc(1, pInfoData->info.bytes);
2,384,862✔
454
        if (pBlockInfo->pks[0].pData == NULL) {
2,375,444✔
455
          return terrno;
×
456
        }
457

458
        pBlockInfo->pks[1].pData = taosMemoryCalloc(1, pInfoData->info.bytes);
2,376,137✔
459
        if (pBlockInfo->pks[1].pData == NULL) {
2,378,909✔
460
          taosMemoryFreeClear(pBlockInfo->pks[0].pData);
×
461
          return terrno;
×
462
        }
463

464
        pBlockInfo->pks[0].nData = pInfoData->info.bytes;
2,378,909✔
465
        pBlockInfo->pks[1].nData = pInfoData->info.bytes;
2,380,988✔
466
      }
467

468
      break;
7,100,513✔
469
    }
470
  }
471

472
  return TSDB_CODE_SUCCESS;
205,893,467✔
473
}
474

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

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

498
    res->translate = true;
118,296✔
499
    res->node.resType = pSColumnNode->node.resType;
118,296✔
500

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

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

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

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

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

559
  return DEAL_RES_CONTINUE;
414,036✔
560
}
561

562
int32_t isQualifiedTable(STableKeyInfo* info, SNode* pTagCond, void* metaHandle, bool* pQualified, SStorageAPI* pAPI) {
59,148✔
563
  int32_t     code = TSDB_CODE_SUCCESS;
59,148✔
564
  SMetaReader mr = {0};
59,148✔
565

566
  pAPI->metaReaderFn.initReader(&mr, metaHandle, META_READER_LOCK, &pAPI->metaFn);
59,148✔
567
  code = pAPI->metaReaderFn.getEntryGetUidCache(&mr, info->uid);
59,148✔
568
  if (TSDB_CODE_SUCCESS != code) {
59,148✔
569
    pAPI->metaReaderFn.clearReader(&mr);
×
570
    *pQualified = false;
×
571

572
    return TSDB_CODE_SUCCESS;
×
573
  }
574

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

591
  SNode* pNew = NULL;
59,148✔
592
  code = scalarCalculateConstants(pTagCondTmp, &pNew);
59,148✔
593
  if (TSDB_CODE_SUCCESS != code) {
59,148✔
594
    terrno = code;
×
595
    nodesDestroyNode(pTagCondTmp);
×
596
    *pQualified = false;
×
597

598
    return code;
×
599
  }
600

601
  SValueNode* pValue = (SValueNode*)pNew;
59,148✔
602
  *pQualified = pValue->datum.b;
59,148✔
603

604
  nodesDestroyNode(pNew);
59,148✔
605
  return TSDB_CODE_SUCCESS;
59,148✔
606
}
607

608
static EDealRes getColumn(SNode** pNode, void* pContext) {
47,184,468✔
609
  tagFilterAssist* pData = (tagFilterAssist*)pContext;
47,184,468✔
610
  SColumnNode*     pSColumnNode = NULL;
47,184,468✔
611
  if (QUERY_NODE_COLUMN == nodeType((*pNode))) {
47,185,845✔
612
    pSColumnNode = *(SColumnNode**)pNode;
15,463,573✔
613
  } else if (QUERY_NODE_FUNCTION == nodeType((*pNode))) {
31,724,552✔
614
    SFunctionNode* pFuncNode = *(SFunctionNode**)(pNode);
685,421✔
615
    if (pFuncNode->funcType == FUNCTION_TYPE_TBNAME) {
685,041✔
616
      pData->code = nodesMakeNode(QUERY_NODE_COLUMN, (SNode**)&pSColumnNode);
634,301✔
617
      if (NULL == pSColumnNode) {
634,301✔
618
        return DEAL_RES_ERROR;
×
619
      }
620
      pSColumnNode->colId = -1;
634,301✔
621
      pSColumnNode->colType = COLUMN_TYPE_TBNAME;
634,301✔
622
      pSColumnNode->node.resType.type = TSDB_DATA_TYPE_VARCHAR;
634,301✔
623
      pSColumnNode->node.resType.bytes = TSDB_TABLE_FNAME_LEN - 1 + VARSTR_HEADER_SIZE;
634,301✔
624
      nodesDestroyNode(*pNode);
634,301✔
625
      *pNode = (SNode*)pSColumnNode;
634,301✔
626
    } else {
627
      return DEAL_RES_CONTINUE;
51,120✔
628
    }
629
  } else {
630
    return DEAL_RES_CONTINUE;
31,037,730✔
631
  }
632

633
  void* data = taosHashGet(pData->colHash, &pSColumnNode->colId, sizeof(pSColumnNode->colId));
16,098,098✔
634
  if (!data) {
16,098,255✔
635
    int32_t tempRes =
636
        taosHashPut(pData->colHash, &pSColumnNode->colId, sizeof(pSColumnNode->colId), pNode, sizeof((*pNode)));
14,782,796✔
637
    if (tempRes != TSDB_CODE_SUCCESS && tempRes != TSDB_CODE_DUP_KEY) {
14,782,178✔
638
      return DEAL_RES_ERROR;
×
639
    }
640
    pSColumnNode->slotId = pData->index++;
14,782,178✔
641
    SColumnInfo cInfo = {.colId = pSColumnNode->colId,
14,780,756✔
642
                         .type = pSColumnNode->node.resType.type,
14,781,770✔
643
                         .bytes = pSColumnNode->node.resType.bytes,
14,780,670✔
644
                         .pk = pSColumnNode->isPk};
14,781,686✔
645
#if TAG_FILTER_DEBUG
646
    qDebug("tagfilter build column info, slotId:%d, colId:%d, type:%d", pSColumnNode->slotId, cInfo.colId, cInfo.type);
647
#endif
648
    void* tmp = taosArrayPush(pData->cInfoList, &cInfo);
14,781,791✔
649
    if (!tmp) {
14,782,727✔
650
      return DEAL_RES_ERROR;
×
651
    }
652
  } else {
653
    SColumnNode* col = *(SColumnNode**)data;
1,315,459✔
654
    pSColumnNode->slotId = col->slotId;
1,315,459✔
655
  }
656

657
  return DEAL_RES_CONTINUE;
16,097,244✔
658
}
659

660
static int32_t createResultData(SDataType* pType, int32_t numOfRows, SScalarParam* pParam) {
13,631,214✔
661
  SColumnInfoData* pColumnData = taosMemoryCalloc(1, sizeof(SColumnInfoData));
13,631,214✔
662
  if (pColumnData == NULL) {
13,632,929✔
663
    return terrno;
×
664
  }
665

666
  pColumnData->info.type = pType->type;
13,632,929✔
667
  pColumnData->info.bytes = pType->bytes;
13,631,531✔
668
  pColumnData->info.scale = pType->scale;
13,630,853✔
669
  pColumnData->info.precision = pType->precision;
13,629,780✔
670

671
  int32_t code = colInfoDataEnsureCapacity(pColumnData, numOfRows, true);
13,629,584✔
672
  if (code != TSDB_CODE_SUCCESS) {
13,627,033✔
673
    terrno = code;
×
674
    releaseColInfoData(pColumnData);
×
675
    return terrno;
×
676
  }
677

678
  pParam->columnData = pColumnData;
13,627,033✔
679
  pParam->colAlloced = true;
13,628,939✔
680
  return TSDB_CODE_SUCCESS;
13,628,797✔
681
}
682

683
static void releaseColInfoData(void* pCol) {
2,192,420✔
684
  if (pCol) {
2,192,420✔
685
    SColumnInfoData* col = (SColumnInfoData*)pCol;
2,193,108✔
686
    colDataDestroy(col);
2,193,108✔
687
    taosMemoryFree(col);
2,193,108✔
688
  }
689
}
2,192,570✔
690

691
void freeItem(void* p) {
189,322,912✔
692
  STUidTagInfo* pInfo = p;
189,322,912✔
693
  if (pInfo->pTagVal != NULL) {
189,322,912✔
694
    taosMemoryFree(pInfo->pTagVal);
189,048,527✔
695
  }
696
}
189,470,151✔
697

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

704
static int compareTagDataEntry(const void* a, const void* b) {
42,834✔
705
  STagDataEntry* p1 = (STagDataEntry*)a;
42,834✔
706
  STagDataEntry* p2 = (STagDataEntry*)b;
42,834✔
707
  return compareInt16Val(&p1->colId, &p2->colId);
42,834✔
708
}
709

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

725
    (void)memcpy(pStart, &entry->colId, sizeof(col_id_t));
42,834✔
726
    pStart += sizeof(col_id_t);
42,834✔
727

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

773
  return TSDB_CODE_SUCCESS;
21,417✔
774
}
775

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

785
  STagDataEntry entry = {0};
42,834✔
786
  entry.colId = pColNode->colId;
42,834✔
787
  entry.pValueNode = (SNode*)pValueNode;
42,834✔
788
  entry.bytes = pColNode->node.resType.bytes;
42,834✔
789
  void* _tmp = taosArrayPush(pIdWithValue, &entry);
42,834✔
790
}
42,834✔
791

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

801
  if (nodeType(pTagCond) == QUERY_NODE_OPERATOR) {
21,417✔
802
    extractTagDataEntry((SOperatorNode*)pTagCond, pIdWithVal);
×
803
  } else if (nodeType(pTagCond) == QUERY_NODE_LOGIC_CONDITION) {
21,417✔
804
    SNode* pChild = NULL;
21,417✔
805
    FOREACH(pChild, ((SLogicConditionNode*)pTagCond)->pParameterList) {
64,251✔
806
      extractTagDataEntry((SOperatorNode*)pChild, pIdWithVal);
42,834✔
807
    }
808
  }
809

810
  taosArraySort(pIdWithVal, compareTagDataEntry);
21,417✔
811

812
  return TSDB_CODE_SUCCESS;
21,417✔
813
}
814

815
static int32_t genStableTagFilterDigest(const SNode* pTagCond, T_MD5_CTX* pContext) {
21,417✔
816
  if (pTagCond == NULL) {
21,417✔
817
    return TSDB_CODE_SUCCESS;
×
818
  }
819

820
  char*   payload = NULL;
21,417✔
821
  int32_t len = 0;
21,417✔
822
  int32_t code = TSDB_CODE_SUCCESS;
21,417✔
823
  int32_t lino = 0;
21,417✔
824

825
  SArray* pIdWithVal = taosArrayInit(TARRAY_MIN_SIZE, sizeof(STagDataEntry));
21,417✔
826
  code = extractTagFilterTagDataEntries(pTagCond, pIdWithVal);
21,417✔
827
  QUERY_CHECK_CODE(code, lino, _end);
21,417✔
828
  for (int32_t i = 0; i < taosArrayGetSize(pIdWithVal); ++i) {
64,251✔
829
    STagDataEntry* pEntry = taosArrayGet(pIdWithVal, i);
42,834✔
830
    len += sizeof(col_id_t) + pEntry->bytes;
42,834✔
831
  }
832
  code = buildTagDataEntryKey(pIdWithVal, &payload, len);
21,417✔
833
  QUERY_CHECK_CODE(code, lino, _end);
21,417✔
834

835
  tMD5Init(pContext);
21,417✔
836
  tMD5Update(pContext, (uint8_t*)payload, (uint32_t)len);
21,417✔
837
  tMD5Final(pContext);
21,417✔
838

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

849
static int32_t genTagFilterDigest(const SNode* pTagCond, T_MD5_CTX* pContext) {
65,811✔
850
  if (pTagCond == NULL) {
65,811✔
851
    return TSDB_CODE_SUCCESS;
62,517✔
852
  }
853

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

862
  tMD5Init(pContext);
3,294✔
863
  tMD5Update(pContext, (uint8_t*)payload, (uint32_t)len);
3,294✔
864
  tMD5Final(pContext);
3,294✔
865

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

872
  taosMemoryFree(payload);
3,294✔
873
  return TSDB_CODE_SUCCESS;
3,294✔
874
}
875

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

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

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

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

903
int32_t qGetColumnsFromNodeList(void* data, bool isList, SArray** pColList) {
13,451,922✔
904
  int32_t code = TSDB_CODE_SUCCESS;
13,451,922✔
905
  tagFilterAssist ctx = {0};
13,451,922✔
906
  ctx.colHash = taosHashInit(4, taosGetDefaultHashFunction(TSDB_DATA_TYPE_SMALLINT), false, HASH_NO_LOCK);
13,452,882✔
907
  if (ctx.colHash == NULL) {
13,452,832✔
908
    code = terrno;
×
909
    goto end;
×
910
  }
911

912
  ctx.index = 0;
13,452,832✔
913
  ctx.cInfoList = taosArrayInit(4, sizeof(SColumnInfo));
13,452,832✔
914
  if (ctx.cInfoList == NULL) {
13,453,240✔
915
    code = terrno;
734✔
916
    goto end;
×
917
  }
918

919
  if (isList) {
13,452,506✔
920
    SNode* pNode = NULL;
2,011,101✔
921
    FOREACH(pNode, (SNodeList*)data) {
4,211,715✔
922
      nodesRewriteExprPostOrder(&pNode, getColumn, (void*)&ctx);
2,200,614✔
923
      if (TSDB_CODE_SUCCESS != ctx.code) {
2,200,614✔
924
        code = ctx.code;
×
925
        goto end;
×
926
      }
927
      REPLACE_NODE(pNode);
2,200,614✔
928
    }
929
  } else {
930
    SNode* pNode = (SNode*)data;
11,441,405✔
931
    nodesRewriteExprPostOrder(&pNode, getColumn, (void*)&ctx);
11,441,405✔
932
    if (TSDB_CODE_SUCCESS != ctx.code) {
11,441,294✔
933
      code = ctx.code;
×
934
      goto end;
×
935
    }
936
  }
937
  
938
  if (pColList != NULL) *pColList = ctx.cInfoList;
13,451,870✔
939
  ctx.cInfoList = NULL;
13,451,218✔
940

941
end:
13,451,743✔
942
  taosHashCleanup(ctx.colHash);
13,452,836✔
943
  taosArrayDestroy(ctx.cInfoList);
13,448,589✔
944
  return code;
13,450,731✔
945
}
946

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

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

1027
  int32_t rows = taosArrayGetSize(pTableListInfo->pTableList);
107,663✔
1028
  if (rows == 0) {
107,663✔
1029
    return;
×
1030
  }
1031

1032
  pUidTagList = taosArrayInit(8, sizeof(STUidTagInfo));
107,663✔
1033
  QUERY_CHECK_NULL(pUidTagList, code, lino, end, terrno);
107,663✔
1034

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

1052
  SArray* pColList = NULL;
107,663✔
1053
  code = qGetColumnsFromNodeList(group, true, &pColList);
107,663✔
1054
  if (code != TSDB_CODE_SUCCESS) {
107,663✔
1055
    goto end;
×
1056
  }
1057

1058
  for (int32_t i = 0; i < taosArrayGetSize(pColList); ++i) {
279,406✔
1059
    SColumnInfo* tmp = (SColumnInfo*)taosArrayGet(pColList, i);
171,743✔
1060
    if (tmp != NULL && tmp->colId == -1) {
171,743✔
1061
      tbNameIndex = i;
107,663✔
1062
    }
1063
  }
1064
  
1065
  int32_t numOfTables = taosArrayGetSize(pUidTagList);
107,283✔
1066
  pResBlock = createTagValBlockForFilter(pColList, numOfTables, pUidTagList, pVnode, pAPI);
107,283✔
1067
  taosArrayDestroy(pColList);
107,663✔
1068
  if (pResBlock == NULL) {
107,663✔
1069
    code = terrno;
×
1070
    goto end;
×
1071
  }
1072

1073
  pBlockList = taosArrayInit(2, POINTER_BYTES);
107,663✔
1074
  QUERY_CHECK_NULL(pBlockList, code, lino, end, terrno);
107,663✔
1075

1076
  void* tmp = taosArrayPush(pBlockList, &pResBlock);
107,663✔
1077
  QUERY_CHECK_NULL(tmp, code, lino, end, terrno);
107,663✔
1078

1079
  groupData = taosArrayInit(2, POINTER_BYTES);
107,663✔
1080
  QUERY_CHECK_NULL(groupData, code, lino, end, terrno);
107,663✔
1081

1082
  SNode* pNode = NULL;
107,663✔
1083
  FOREACH(pNode, group) {
281,686✔
1084
    SScalarParam output = {0};
174,023✔
1085

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

1100
      default:
×
1101
        code = TSDB_CODE_OPS_NOT_SUPPORT;
×
1102
        goto end;
×
1103
    }
1104

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

1118
    if (code != TSDB_CODE_SUCCESS) {
174,023✔
1119
      releaseColInfoData(output.columnData);
×
1120
      goto end;
×
1121
    }
1122

1123
    void* tmp = taosArrayPush(groupData, &output.columnData);
174,023✔
1124
    QUERY_CHECK_NULL(tmp, code, lino, end, terrno);
174,023✔
1125
  }
1126

1127
  for (int i = 0; i < rows; i++) {
453,713✔
1128
    gInfo = taosArrayInit(taosArrayGetSize(groupData), sizeof(SStreamGroupValue));
346,050✔
1129
    QUERY_CHECK_NULL(gInfo, code, lino, end, terrno);
346,050✔
1130

1131
    STableKeyInfo* info = taosArrayGet(pTableListInfo->pTableList, i);
346,050✔
1132
    QUERY_CHECK_NULL(info, code, lino, end, terrno);
346,050✔
1133

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

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

1159
end:
107,663✔
1160
  blockDataDestroy(pResBlock);
107,663✔
1161
  taosArrayDestroy(pBlockList);
107,663✔
1162
  taosArrayDestroyEx(pUidTagList, freeItem);
107,663✔
1163
  taosArrayDestroyP(groupData, releaseColInfoData);
107,663✔
1164
  taosArrayDestroyEx(gInfo, tDestroySStreamGroupValue);
107,663✔
1165

1166
  if (code != TSDB_CODE_SUCCESS) {
107,663✔
1167
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
1168
  }
1169
}
1170

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

1183
  int32_t rows = taosArrayGetSize(pTableListInfo->pTableList);
1,903,438✔
1184
  if (rows == 0) {
1,903,438✔
1185
    return TSDB_CODE_SUCCESS;
×
1186
  } 
1187

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

1199
    nodesFree(listNode);
×
1200

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

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

1214
  pUidTagList = taosArrayInit(8, sizeof(STUidTagInfo));
1,903,438✔
1215
  QUERY_CHECK_NULL(pUidTagList, code, lino, end, terrno);
1,903,438✔
1216

1217
  for (int32_t i = 0; i < rows; ++i) {
11,090,721✔
1218
    STableKeyInfo* pkeyInfo = taosArrayGet(pTableListInfo->pTableList, i);
9,187,175✔
1219
    QUERY_CHECK_NULL(pkeyInfo, code, lino, end, terrno);
9,187,175✔
1220
    STUidTagInfo info = {.uid = pkeyInfo->uid};
9,187,175✔
1221
    void*        tmp = taosArrayPush(pUidTagList, &info);
9,187,283✔
1222
    QUERY_CHECK_NULL(tmp, code, lino, end, terrno);
9,187,283✔
1223
  }
1224

1225
  if (taosArrayGetSize(pUidTagList) > 0) {
1,903,546✔
1226
    code = pAPI->metaFn.getTableTagsByUid(pVnode, pTableListInfo->idInfo.suid, pUidTagList);
1,903,438✔
1227
  } else {
1228
    code = pAPI->metaFn.getTableTags(pVnode, pTableListInfo->idInfo.suid, pUidTagList);
×
1229
  }
1230
  if (code != TSDB_CODE_SUCCESS) {
1,903,288✔
1231
    goto end;
×
1232
  }
1233

1234
  SArray* pColList = NULL;
1,903,288✔
1235
  code = qGetColumnsFromNodeList(group, true, &pColList); 
1,903,288✔
1236
  if (code != TSDB_CODE_SUCCESS) {
1,903,288✔
1237
    goto end;
×
1238
  }
1239

1240
  int32_t numOfTables = taosArrayGetSize(pUidTagList);
1,903,288✔
1241
  pResBlock = createTagValBlockForFilter(pColList, numOfTables, pUidTagList, pVnode, pAPI);
1,903,288✔
1242
  taosArrayDestroy(pColList);
1,901,987✔
1243
  if (pResBlock == NULL) {
1,902,600✔
1244
    code = terrno;
×
1245
    goto end;
×
1246
  }
1247

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

1251
  pBlockList = taosArrayInit(2, POINTER_BYTES);
1,902,600✔
1252
  QUERY_CHECK_NULL(pBlockList, code, lino, end, terrno);
1,903,438✔
1253

1254
  void* tmp = taosArrayPush(pBlockList, &pResBlock);
1,902,600✔
1255
  QUERY_CHECK_NULL(tmp, code, lino, end, terrno);
1,902,600✔
1256

1257
  groupData = taosArrayInit(2, POINTER_BYTES);
1,902,600✔
1258
  QUERY_CHECK_NULL(groupData, code, lino, end, terrno);
1,902,909✔
1259

1260
  SNode* pNode = NULL;
1,902,221✔
1261
  FOREACH(pNode, group) {
3,928,507✔
1262
    SScalarParam output = {0};
2,025,903✔
1263

1264
    switch (nodeType(pNode)) {
2,025,374✔
1265
      case QUERY_NODE_VALUE:
×
1266
        break;
×
1267
      case QUERY_NODE_COLUMN:
2,018,397✔
1268
      case QUERY_NODE_OPERATOR:
1269
      case QUERY_NODE_FUNCTION: {
1270
        SExprNode* expNode = (SExprNode*)pNode;
2,018,397✔
1271
        code = createResultData(&expNode->resType, rows, &output);
2,018,397✔
1272
        if (code != TSDB_CODE_SUCCESS) {
2,017,600✔
1273
          goto end;
×
1274
        }
1275
        break;
2,017,600✔
1276
      }
1277
      case QUERY_NODE_REMOTE_VALUE: {
7,356✔
1278
        SRemoteValueNode* pRemote = (SRemoteValueNode*)pNode;
7,356✔
1279
        code = qFetchRemoteValue(gTaskScalarExtra.pSubJobCtx, pRemote->subQIdx, pRemote);
7,356✔
1280
        QUERY_CHECK_CODE(code, lino, end);
7,356✔
1281
        break;
7,356✔
1282
      }
1283
      
1284
      default:
×
1285
        code = TSDB_CODE_OPS_NOT_SUPPORT;
×
1286
        goto end;
×
1287
    }
1288

1289
    if (nodeType(pNode) == QUERY_NODE_COLUMN) {
2,024,956✔
1290
      SColumnNode*     pSColumnNode = (SColumnNode*)pNode;
2,000,152✔
1291
      SColumnInfoData* pColInfo = (SColumnInfoData*)taosArrayGet(pResBlock->pDataBlock, pSColumnNode->slotId);
2,000,152✔
1292
      QUERY_CHECK_NULL(pColInfo, code, lino, end, terrno);
1,999,314✔
1293
      code = colDataAssign(output.columnData, pColInfo, rows, NULL);
1,999,314✔
1294
    } else if (nodeType(pNode) == QUERY_NODE_VALUE) {
24,041✔
1295
      continue;
7,356✔
1296
    } else {
1297
      gTaskScalarExtra.pStreamInfo = NULL;
17,598✔
1298
      gTaskScalarExtra.pStreamRange = NULL;
17,598✔
1299
      code = scalarCalculate(pNode, pBlockList, &output, &gTaskScalarExtra);
17,598✔
1300
    }
1301

1302
    if (code != TSDB_CODE_SUCCESS) {
2,017,076✔
1303
      releaseColInfoData(output.columnData);
×
1304
      goto end;
×
1305
    }
1306

1307
    void* tmp = taosArrayPush(groupData, &output.columnData);
2,019,080✔
1308
    QUERY_CHECK_NULL(tmp, code, lino, end, terrno);
2,019,080✔
1309
  }
1310

1311
  int32_t keyLen = 0;
1,902,445✔
1312
  SNode*  node;
1313
  FOREACH(node, group) {
3,927,085✔
1314
    SExprNode* pExpr = (SExprNode*)node;
2,026,286✔
1315
    keyLen += pExpr->resType.bytes;
2,026,286✔
1316
  }
1317

1318
  int32_t nullFlagSize = sizeof(int8_t) * LIST_LENGTH(group);
1,902,325✔
1319
  keyLen += nullFlagSize;
1,902,175✔
1320

1321
  keyBuf = taosMemoryCalloc(1, keyLen);
1,902,175✔
1322
  if (keyBuf == NULL) {
1,902,445✔
1323
    code = terrno;
×
1324
    goto end;
×
1325
  }
1326

1327
  if (initRemainGroups) {
1,902,445✔
1328
    pTableListInfo->remainGroups =
856,031✔
1329
        taosHashInit(rows, taosGetDefaultHashFunction(TSDB_DATA_TYPE_BIGINT), false, HASH_NO_LOCK);
855,876✔
1330
    if (pTableListInfo->remainGroups == NULL) {
856,031✔
1331
      code = terrno;
×
1332
      goto end;
×
1333
    }
1334
  }
1335

1336
  for (int i = 0; i < rows; i++) {
11,086,355✔
1337
    STableKeyInfo* info = taosArrayGet(pTableListInfo->pTableList, i);
9,182,520✔
1338
    QUERY_CHECK_NULL(info, code, lino, end, terrno);
9,185,246✔
1339

1340
    if (groupIdMap != NULL){
9,185,246✔
1341
      gInfo = taosArrayInit(taosArrayGetSize(groupData), sizeof(SStreamGroupValue));
192,105✔
1342
    }
1343
    
1344
    char* isNull = (char*)keyBuf;
9,183,788✔
1345
    char* pStart = (char*)keyBuf + sizeof(int8_t) * LIST_LENGTH(group);
9,183,788✔
1346
    for (int j = 0; j < taosArrayGetSize(groupData); j++) {
19,010,373✔
1347
      SColumnInfoData* pValue = (SColumnInfoData*)taosArrayGetP(groupData, j);
9,826,626✔
1348

1349
      if (groupIdMap != NULL && gInfo != NULL) {
9,826,810✔
1350
        int32_t ret = buildGroupInfo(pValue, i, gInfo);
215,801✔
1351
        if (ret != TSDB_CODE_SUCCESS) {
215,801✔
1352
          qError("buildGroupInfo failed at line %d since %s", __LINE__, tstrerror(ret));
×
1353
          taosArrayDestroyEx(gInfo, tDestroySStreamGroupValue);
×
1354
          gInfo = NULL;
×
1355
        }
1356
      }
1357
      
1358
      if (colDataIsNull_s(pValue, i)) {
19,652,632✔
1359
        isNull[j] = 1;
94,961✔
1360
      } else {
1361
        isNull[j] = 0;
9,730,861✔
1362
        char* data = colDataGetData(pValue, i);
9,730,173✔
1363
        if (pValue->info.type == TSDB_DATA_TYPE_JSON) {
9,731,240✔
1364
          // if (tTagIsJson(data)) {
1365
          //   code = TSDB_CODE_QRY_JSON_IN_GROUP_ERROR;
1366
          //   goto end;
1367
          // }
1368
          if (tTagIsJsonNull(data)) {
89,650✔
1369
            isNull[j] = 1;
×
1370
            continue;
×
1371
          }
1372
          int32_t len = getJsonValueLen(data);
89,650✔
1373
          memcpy(pStart, data, len);
89,650✔
1374
          pStart += len;
89,650✔
1375
        } else if (IS_VAR_DATA_TYPE(pValue->info.type)) {
9,640,907✔
1376
          if (IS_STR_DATA_BLOB(pValue->info.type)) {
6,583,896✔
1377
            if (blobDataTLen(data) > TSDB_MAX_BLOB_LEN) {
×
1378
              code = TSDB_CODE_TDB_INVALID_TABLE_SCHEMA_VER;
×
1379
              goto end;
×
1380
            }
1381
            memcpy(pStart, data, blobDataTLen(data));
×
1382
            pStart += blobDataTLen(data);
×
1383
          } else {
1384
            if (varDataTLen(data) > pValue->info.bytes) {
6,584,230✔
1385
              code = TSDB_CODE_TDB_INVALID_TABLE_SCHEMA_VER;
×
1386
              goto end;
×
1387
            }
1388
            memcpy(pStart, data, varDataTLen(data));
6,583,625✔
1389
            pStart += varDataTLen(data);
6,583,625✔
1390
          }
1391
        } else {
1392
          memcpy(pStart, data, pValue->info.bytes);
3,057,416✔
1393
          pStart += pValue->info.bytes;
3,057,566✔
1394
        }
1395
      }
1396
    }
1397

1398
    int32_t len = (int32_t)(pStart - (char*)keyBuf);
9,182,026✔
1399
    info->groupId = calcGroupId(keyBuf, len);
9,182,026✔
1400
    if (groupIdMap != NULL && gInfo != NULL) {
9,184,327✔
1401
      int32_t ret = taosHashPut(groupIdMap, &info->groupId, sizeof(info->groupId), &gInfo, POINTER_BYTES);
191,742✔
1402
      if (ret != TSDB_CODE_SUCCESS) {
192,105✔
1403
        qError("put groupid to map failed at line %d since %s", __LINE__, tstrerror(ret));
×
1404
        taosArrayDestroyEx(gInfo, tDestroySStreamGroupValue);
×
1405
      }
1406
      qDebug("put groupid to map gid:%" PRIu64, info->groupId);
192,105✔
1407
      gInfo = NULL;
191,742✔
1408
    }
1409
    if (initRemainGroups) {
9,184,327✔
1410
      // groupId ~ table uid
1411
      code = taosHashPut(pTableListInfo->remainGroups, &(info->groupId), sizeof(info->groupId), &(info->uid),
4,580,520✔
1412
                         sizeof(info->uid));
1413
      if (code == TSDB_CODE_DUP_KEY) {
4,580,098✔
1414
        code = TSDB_CODE_SUCCESS;
843,155✔
1415
      }
1416
      QUERY_CHECK_CODE(code, lino, end);
4,580,098✔
1417
    }
1418
  }
1419

1420
  if (tsTagFilterCache && groupIdMap == NULL) {
1,903,835✔
1421
    tableList = taosArrayDup(pTableListInfo->pTableList, NULL);
×
1422
    QUERY_CHECK_NULL(tableList, code, lino, end, terrno);
×
1423

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

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

1433
end:
1,902,997✔
1434
  taosMemoryFreeClear(keyBuf);
1,902,600✔
1435
  blockDataDestroy(pResBlock);
1,903,438✔
1436
  taosArrayDestroy(pBlockList);
1,903,288✔
1437
  taosArrayDestroyEx(pUidTagList, freeItem);
1,903,288✔
1438
  taosArrayDestroyP(groupData, releaseColInfoData);
1,903,133✔
1439
  taosArrayDestroyEx(gInfo, tDestroySStreamGroupValue);
1,903,288✔
1440

1441
  if (code != TSDB_CODE_SUCCESS) {
1,902,600✔
1442
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
1443
  }
1444
  return code;
1,902,600✔
1445
}
1446

1447
static int32_t nameComparFn(const void* p1, const void* p2) {
723,464✔
1448
  const char* pName1 = *(const char**)p1;
723,464✔
1449
  const char* pName2 = *(const char**)p2;
723,464✔
1450

1451
  int32_t ret = strcmp(pName1, pName2);
723,464✔
1452
  if (ret == 0) {
723,464✔
1453
    return 0;
18,288✔
1454
  } else {
1455
    return (ret > 0) ? 1 : -1;
705,176✔
1456
  }
1457
}
1458

1459
static SArray* getTableNameList(const SNodeListNode* pList) {
423,204✔
1460
  int32_t    code = TSDB_CODE_SUCCESS;
423,204✔
1461
  int32_t    lino = 0;
423,204✔
1462
  int32_t    len = LIST_LENGTH(pList->pNodeList);
423,204✔
1463
  SListCell* cell = pList->pNodeList->pHead;
423,204✔
1464

1465
  SArray* pTbList = taosArrayInit(len, POINTER_BYTES);
423,009✔
1466
  QUERY_CHECK_NULL(pTbList, code, lino, _end, terrno);
423,204✔
1467

1468
  for (int i = 0; i < pList->pNodeList->length; i++) {
1,149,454✔
1469
    SValueNode* valueNode = (SValueNode*)cell->pNode;
726,250✔
1470
    if (!IS_VAR_DATA_TYPE(valueNode->node.resType.type)) {
726,250✔
1471
      terrno = TSDB_CODE_INVALID_PARA;
×
1472
      taosArrayDestroy(pTbList);
×
1473
      return NULL;
×
1474
    }
1475

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

1482
  size_t numOfTables = taosArrayGetSize(pTbList);
423,204✔
1483

1484
  // order the name
1485
  taosArraySort(pTbList, nameComparFn);
423,204✔
1486

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

1495
  for (int32_t i = 1; i < numOfTables; ++i) {
726,250✔
1496
    char** name = taosArrayGetLast(pNewList);
303,046✔
1497
    char** nameInOldList = taosArrayGet(pTbList, i);
303,046✔
1498
    QUERY_CHECK_NULL(nameInOldList, code, lino, _end, terrno);
303,046✔
1499
    if (strcmp(*name, *nameInOldList) == 0) {
303,046✔
1500
      continue;
9,836✔
1501
    }
1502

1503
    tmp = taosArrayPush(pNewList, nameInOldList);
293,210✔
1504
    QUERY_CHECK_NULL(tmp, code, lino, _end, terrno);
293,210✔
1505
  }
1506

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

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

1520
  if (u1 == u2) {
×
1521
    return 0;
×
1522
  }
1523

1524
  return u1 < u2 ? -1 : 1;
×
1525
}
1526

1527
static int32_t filterTableInfoCompare(const void* a, const void* b) {
18,436,562✔
1528
  STUidTagInfo* p1 = (STUidTagInfo*)a;
18,436,562✔
1529
  STUidTagInfo* p2 = (STUidTagInfo*)b;
18,436,562✔
1530

1531
  if (p1->uid == p2->uid) {
18,436,562✔
1532
    return 0;
×
1533
  }
1534

1535
  return p1->uid < p2->uid ? -1 : 1;
18,436,562✔
1536
}
1537

1538
static FilterCondType checkTagCond(SNode* cond) {
13,335,197✔
1539
  if (nodeType(cond) == QUERY_NODE_OPERATOR) {
13,335,197✔
1540
    return FILTER_NO_LOGIC;
11,277,748✔
1541
  }
1542
  if (nodeType(cond) != QUERY_NODE_LOGIC_CONDITION || ((SLogicConditionNode*)cond)->condType != LOGIC_COND_TYPE_AND) {
2,058,968✔
1543
    return FILTER_AND;
216,392✔
1544
  }
1545
  return FILTER_OTHER;
1,841,776✔
1546
}
1547

1548
static int32_t optimizeTbnameInCond(void* pVnode, int64_t suid, SArray* list, SNode* cond, SStorageAPI* pAPI) {
13,335,714✔
1549
  int32_t ret = -1;
13,335,714✔
1550
  int32_t ntype = nodeType(cond);
13,335,714✔
1551

1552
  if (ntype == QUERY_NODE_OPERATOR) {
13,335,914✔
1553
    ret = optimizeTbnameInCondImpl(pVnode, list, cond, pAPI, suid);
11,277,346✔
1554
  }
1555

1556
  if (ntype != QUERY_NODE_LOGIC_CONDITION || ((SLogicConditionNode*)cond)->condType != LOGIC_COND_TYPE_AND) {
13,335,914✔
1557
    return ret;
11,493,536✔
1558
  }
1559

1560
  bool                 hasTbnameCond = false;
1,842,176✔
1561
  SLogicConditionNode* pNode = (SLogicConditionNode*)cond;
1,842,176✔
1562
  SNodeList*           pList = (SNodeList*)pNode->pParameterList;
1,842,176✔
1563

1564
  int32_t len = LIST_LENGTH(pList);
1,842,384✔
1565
  if (len <= 0) {
1,842,376✔
1566
    return ret;
×
1567
  }
1568

1569
  SListCell* cell = pList->pHead;
1,842,376✔
1570
  for (int i = 0; i < len; i++) {
5,949,522✔
1571
    if (cell == NULL) break;
4,113,378✔
1572
    if (optimizeTbnameInCondImpl(pVnode, list, cell->pNode, pAPI, suid) == 0) {
4,113,378✔
1573
      hasTbnameCond = true;
6,230✔
1574
      break;
6,230✔
1575
    }
1576
    cell = cell->pNext;
4,107,556✔
1577
  }
1578

1579
  taosArraySort(list, filterTableInfoCompare);
1,842,374✔
1580
  taosArrayRemoveDuplicate(list, filterTableInfoCompare, NULL);
1,842,384✔
1581

1582
  if (hasTbnameCond) {
1,842,384✔
1583
    ret = pAPI->metaFn.getTableTagsByUid(pVnode, suid, list);
6,230✔
1584
  }
1585

1586
  return ret;
1,842,176✔
1587
}
1588

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

1596
  SOperatorNode* pNode = (SOperatorNode*)pTagCond;
15,385,584✔
1597
  if (pNode->opType != OP_TYPE_IN) {
15,385,584✔
1598
    return -1;
14,655,736✔
1599
  }
1600

1601
  if ((pNode->pLeft != NULL && ((nodeType(pNode->pLeft) == QUERY_NODE_FUNCTION &&
729,232✔
1602
                                 ((SFunctionNode*)pNode->pLeft)->funcType == FUNCTION_TYPE_TBNAME)) ||
423,204✔
1603
       (nodeType(pNode->pLeft) == QUERY_NODE_COLUMN && ((SColumnNode*)pNode->pLeft)->colType == COLUMN_TYPE_TBNAME)) &&
306,028✔
1604
      (pNode->pRight != NULL && nodeType(pNode->pRight) == QUERY_NODE_NODE_LIST)) {
423,204✔
1605
    SNodeListNode* pList = (SNodeListNode*)pNode->pRight;
423,204✔
1606

1607
    int32_t len = LIST_LENGTH(pList->pNodeList);
423,204✔
1608
    if (len <= 0) {
423,204✔
1609
      return -1;
×
1610
    }
1611

1612
    SArray*   pTbList = getTableNameList(pList);
423,204✔
1613
    int32_t   numOfTables = taosArrayGetSize(pTbList);
423,204✔
1614
    SHashObj* uHash = NULL;
423,204✔
1615

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

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

1638
    for (int i = 0; i < numOfTables; i++) {
1,123,258✔
1639
      char* name = taosArrayGetP(pTbList, i);
706,578✔
1640

1641
      uint64_t uid = 0, csuid = 0;
706,578✔
1642
      if (pStoreAPI->metaFn.getTableUidByName(pVnode, name, &uid) == 0) {
706,578✔
1643
        ETableType tbType = TSDB_TABLE_MAX;
423,320✔
1644
        if (pStoreAPI->metaFn.getTableTypeSuidByName(pVnode, name, &tbType, &csuid) == 0 &&
423,320✔
1645
            tbType == TSDB_CHILD_TABLE) {
423,320✔
1646
          if (suid != csuid) {
416,796✔
1647
            continue;
856✔
1648
          }
1649
          if (NULL == uHash || taosHashGet(uHash, &uid, sizeof(uid)) == NULL) {
415,940✔
1650
            STUidTagInfo s = {.uid = uid, .name = name, .pTagVal = NULL};
414,762✔
1651
            void*        tmp = taosArrayPush(pExistedUidList, &s);
414,762✔
1652
            if (!tmp) {
414,762✔
1653
              return terrno;
×
1654
            }
1655
          }
1656
        } else {
1657
          taosArrayDestroy(pTbList);
6,524✔
1658
          taosHashCleanup(uHash);
6,524✔
1659
          return -1;
6,524✔
1660
        }
1661
      } else {
1662
        //        qWarn("failed to get tableIds from by table name: %s, reason: %s", name, tstrerror(terrno));
1663
        terrno = 0;
283,258✔
1664
      }
1665
    }
1666

1667
    taosHashCleanup(uHash);
416,680✔
1668
    taosArrayDestroy(pTbList);
416,680✔
1669
    return 0;
416,680✔
1670
  }
1671

1672
  return -1;
306,028✔
1673
}
1674

1675
SSDataBlock* createTagValBlockForFilter(SArray* pColList, int32_t numOfTables, SArray* pUidTagList, void* pVnode,
14,114,001✔
1676
                                        SStorageAPI* pStorageAPI) {
1677
  int32_t      code = TSDB_CODE_SUCCESS;
14,114,001✔
1678
  int32_t      lino = 0;
14,114,001✔
1679
  SSDataBlock* pResBlock = NULL;
14,114,001✔
1680
  code = createDataBlock(&pResBlock);
14,114,650✔
1681
  QUERY_CHECK_CODE(code, lino, _end);
14,114,276✔
1682

1683
  for (int32_t i = 0; i < taosArrayGetSize(pColList); ++i) {
29,560,300✔
1684
    SColumnInfoData colInfo = {0};
15,443,339✔
1685
    void*           tmp = taosArrayGet(pColList, i);
15,443,339✔
1686
    QUERY_CHECK_NULL(tmp, code, lino, _end, terrno);
15,444,713✔
1687
    colInfo.info = *(SColumnInfo*)tmp;
15,444,713✔
1688
    code = blockDataAppendColInfo(pResBlock, &colInfo);
15,443,911✔
1689
    QUERY_CHECK_CODE(code, lino, _end);
15,445,559✔
1690
  }
1691

1692
  code = blockDataEnsureCapacity(pResBlock, numOfTables);
14,114,466✔
1693
  if (code != TSDB_CODE_SUCCESS) {
14,113,757✔
1694
    terrno = code;
×
1695
    blockDataDestroy(pResBlock);
×
1696
    return NULL;
×
1697
  }
1698

1699
  pResBlock->info.rows = numOfTables;
14,113,757✔
1700

1701
  int32_t numOfCols = taosArrayGetSize(pResBlock->pDataBlock);
14,114,192✔
1702

1703
  for (int32_t i = 0; i < numOfTables; i++) {
204,532,861✔
1704
    STUidTagInfo* p1 = taosArrayGet(pUidTagList, i);
190,415,551✔
1705
    QUERY_CHECK_NULL(p1, code, lino, _end, terrno);
190,378,847✔
1706

1707
    for (int32_t j = 0; j < numOfCols; j++) {
388,712,480✔
1708
      SColumnInfoData* pColInfo = (SColumnInfoData*)taosArrayGet(pResBlock->pDataBlock, j);
198,289,198✔
1709
      QUERY_CHECK_NULL(pColInfo, code, lino, _end, terrno);
198,276,129✔
1710

1711
      if (pColInfo->info.colId == -1) {  // tbname
198,276,129✔
1712
        char str[TSDB_TABLE_FNAME_LEN + VARSTR_HEADER_SIZE] = {0};
8,033,439✔
1713
        if (p1->name != NULL) {
8,030,698✔
1714
          STR_TO_VARSTR(str, p1->name);
414,762✔
1715
        } else {  // name is not retrieved during filter
1716
          code = pStorageAPI->metaFn.getTableNameByUid(pVnode, p1->uid, str);
7,620,060✔
1717
          QUERY_CHECK_CODE(code, lino, _end);
7,621,818✔
1718
        }
1719

1720
        code = colDataSetVal(pColInfo, i, str, false);
8,036,385✔
1721
        QUERY_CHECK_CODE(code, lino, _end);
8,038,344✔
1722
#if TAG_FILTER_DEBUG
1723
        qDebug("tagfilter uid:%ld, tbname:%s", *uid, str + 2);
1724
#endif
1725
      } else {
1726
        STagVal tagVal = {0};
190,245,780✔
1727
        tagVal.cid = pColInfo->info.colId;
190,213,135✔
1728
        if (p1->pTagVal == NULL) {
190,223,519✔
1729
          colDataSetNULL(pColInfo, i);
9,075✔
1730
        } else {
1731
          const char* p = pStorageAPI->metaFn.extractTagVal(p1->pTagVal, pColInfo->info.type, &tagVal);
190,240,080✔
1732

1733
          if (p == NULL || (pColInfo->info.type == TSDB_DATA_TYPE_JSON && ((STag*)p)->nTag == 0)) {
190,256,611✔
1734
            colDataSetNULL(pColInfo, i);
4,067,193✔
1735
          } else if (pColInfo->info.type == TSDB_DATA_TYPE_JSON) {
186,195,769✔
1736
            code = colDataSetVal(pColInfo, i, p, false);
709,114✔
1737
            QUERY_CHECK_CODE(code, lino, _end);
709,114✔
1738
          } else if (IS_VAR_DATA_TYPE(pColInfo->info.type)) {
299,950,610✔
1739
            if (IS_STR_DATA_BLOB(pColInfo->info.type)) {
114,466,848✔
1740
              QUERY_CHECK_CODE(code = TSDB_CODE_BLOB_NOT_SUPPORT_TAG, lino, _end);
×
1741
            }
1742
            char* tmp = taosMemoryMalloc(tagVal.nData + VARSTR_HEADER_SIZE + 1);
114,469,281✔
1743
            QUERY_CHECK_NULL(tmp, code, lino, _end, terrno);
114,486,777✔
1744
            varDataSetLen(tmp, tagVal.nData);
114,486,777✔
1745
            memcpy(tmp + VARSTR_HEADER_SIZE, tagVal.pData, tagVal.nData);
114,470,885✔
1746
            code = colDataSetVal(pColInfo, i, tmp, false);
114,470,985✔
1747
#if TAG_FILTER_DEBUG
1748
            qDebug("tagfilter varch:%s", tmp + 2);
1749
#endif
1750
            taosMemoryFree(tmp);
114,459,343✔
1751
            QUERY_CHECK_CODE(code, lino, _end);
114,474,722✔
1752
          } else {
1753
            code = colDataSetVal(pColInfo, i, (const char*)&tagVal.i64, false);
71,016,995✔
1754
            QUERY_CHECK_CODE(code, lino, _end);
71,022,227✔
1755
#if TAG_FILTER_DEBUG
1756
            if (pColInfo->info.type == TSDB_DATA_TYPE_INT) {
1757
              qDebug("tagfilter int:%d", *(int*)(&tagVal.i64));
1758
            } else if (pColInfo->info.type == TSDB_DATA_TYPE_DOUBLE) {
1759
              qDebug("tagfilter double:%f", *(double*)(&tagVal.i64));
1760
            }
1761
#endif
1762
          }
1763
        }
1764
      }
1765
    }
1766
  }
1767

1768
_end:
14,120,193✔
1769
  if (code != TSDB_CODE_SUCCESS) {
14,116,433✔
1770
    blockDataDestroy(pResBlock);
3,607✔
1771
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
1772
    terrno = code;
×
1773
    return NULL;
×
1774
  }
1775
  return pResBlock;
14,112,826✔
1776
}
1777

1778
static int32_t doSetQualifiedUid(STableListInfo* pListInfo, SArray* pUidList, const SArray* pUidTagList,
11,430,699✔
1779
                                 bool* pResultList, bool addUid) {
1780
  taosArrayClear(pUidList);
11,430,699✔
1781

1782
  STableKeyInfo info = {.uid = 0, .groupId = 0};
11,431,525✔
1783
  int32_t       numOfTables = taosArrayGetSize(pUidTagList);
11,432,824✔
1784
  for (int32_t i = 0; i < numOfTables; ++i) {
191,354,996✔
1785
    if (pResultList[i]) {
179,914,755✔
1786
      STUidTagInfo* tmpTag = (STUidTagInfo*)taosArrayGet(pUidTagList, i);
78,908,127✔
1787
      if (!tmpTag) {
78,908,399✔
1788
        qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
1789
        return terrno;
×
1790
      }
1791
      uint64_t uid = tmpTag->uid;
78,908,399✔
1792
      qDebug("tagfilter get uid:%" PRId64 ", res:%d", uid, pResultList[i]);
78,899,873✔
1793

1794
      info.uid = uid;
78,904,064✔
1795
      //qInfo("doSetQualifiedUid row:%d added to pTableList", i);
1796
      void* p = taosArrayPush(pListInfo->pTableList, &info);
78,904,064✔
1797
      if (p == NULL) {
78,915,571✔
1798
        return terrno;
×
1799
      }
1800

1801
      if (addUid) {
78,915,571✔
1802
        //qInfo("doSetQualifiedUid row:%d added to pUidList", i);
1803
        void* tmp = taosArrayPush(pUidList, &uid);
20,449✔
1804
        if (tmp == NULL) {
20,449✔
1805
          return terrno;
×
1806
        }
1807
      }
1808
    } else {
1809
      //qInfo("doSetQualifiedUid row:%d failed", i);
1810
    }
1811
  }
1812

1813
  return TSDB_CODE_SUCCESS;
11,440,241✔
1814
}
1815

1816
static int32_t copyExistedUids(SArray* pUidTagList, const SArray* pUidList) {
13,336,716✔
1817
  int32_t code = TSDB_CODE_SUCCESS;
13,336,716✔
1818
  int32_t numOfExisted = taosArrayGetSize(pUidList);
13,336,716✔
1819
  if (numOfExisted == 0) {
13,336,716✔
1820
    return code;
10,408,246✔
1821
  }
1822

1823
  for (int32_t i = 0; i < numOfExisted; ++i) {
35,724,639✔
1824
    uint64_t* uid = taosArrayGet(pUidList, i);
32,796,169✔
1825
    if (!uid) {
32,796,169✔
1826
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
1827
      return terrno;
×
1828
    }
1829
    STUidTagInfo info = {.uid = *uid};
32,796,169✔
1830
    void*        tmp = taosArrayPush(pUidTagList, &info);
32,796,169✔
1831
    if (!tmp) {
32,796,169✔
1832
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
1833
      return code;
×
1834
    }
1835
  }
1836
  return code;
2,928,470✔
1837
}
1838

1839
int32_t doFilterByTagCond(STableListInfo* pListInfo, SArray* pUidList, SNode* pTagCond, void* pVnode,
214,269,060✔
1840
                                 SIdxFltStatus status, SStorageAPI* pAPI, bool addUid, bool* listAdded, void* pStreamInfo) {
1841
  *listAdded = false;
214,269,060✔
1842
  if (pTagCond == NULL) {
214,297,927✔
1843
    return TSDB_CODE_SUCCESS;
200,940,183✔
1844
  }
1845

1846
  terrno = TSDB_CODE_SUCCESS;
13,357,744✔
1847

1848
  int32_t      lino = 0;
13,335,912✔
1849
  int32_t      code = TSDB_CODE_SUCCESS;
13,335,912✔
1850
  SArray*      pBlockList = NULL;
13,335,912✔
1851
  SSDataBlock* pResBlock = NULL;
13,335,912✔
1852
  SScalarParam output = {0};
13,335,510✔
1853
  SArray*      pUidTagList = NULL;
13,335,504✔
1854

1855
  SDataType type = {.type = TSDB_DATA_TYPE_BOOL, .bytes = sizeof(bool)};
13,335,504✔
1856

1857
  //  int64_t stt = taosGetTimestampUs();
1858
  pUidTagList = taosArrayInit(10, sizeof(STUidTagInfo));
13,335,914✔
1859
  QUERY_CHECK_NULL(pUidTagList, code, lino, end, terrno);
13,336,314✔
1860

1861
  code = copyExistedUids(pUidTagList, pUidList);
13,336,314✔
1862
  QUERY_CHECK_CODE(code, lino, end);
13,335,512✔
1863

1864
  FilterCondType condType = checkTagCond(pTagCond);
13,335,512✔
1865

1866
  int32_t filter = optimizeTbnameInCond(pVnode, pListInfo->idInfo.suid, pUidTagList, pTagCond, pAPI);
13,335,316✔
1867
  if (filter == 0) {  // tbname in filter is activated, do nothing and return
13,335,310✔
1868
    taosArrayClear(pUidList);
416,680✔
1869

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

1874
    for (int32_t i = 0; i < numOfRows; ++i) {
3,183,908✔
1875
      STUidTagInfo* pInfo = taosArrayGet(pUidTagList, i);
2,767,228✔
1876
      QUERY_CHECK_NULL(pInfo, code, lino, end, terrno);
2,767,228✔
1877
      void* tmp = taosArrayPush(pUidList, &pInfo->uid);
2,767,228✔
1878
      QUERY_CHECK_NULL(tmp, code, lino, end, terrno);
2,767,228✔
1879
    }
1880
    terrno = 0;
416,680✔
1881
  } else {
1882
    qDebug("pUidTagList size:%d", (int32_t)taosArrayGetSize(pUidTagList));
12,918,630✔
1883
    
1884
    if (((condType == FILTER_NO_LOGIC || condType == FILTER_AND) && status != SFLT_NOT_INDEX) ||
22,651,063✔
1885
          taosArrayGetSize(pUidTagList) > 0) {
9,732,635✔
1886
      code = pAPI->metaFn.getTableTagsByUid(pVnode, pListInfo->idInfo.suid, pUidTagList);
3,685,484✔
1887
    } else {
1888
      code = pAPI->metaFn.getTableTags(pVnode, pListInfo->idInfo.suid, pUidTagList);
9,232,944✔
1889
    }
1890
    if (code != TSDB_CODE_SUCCESS) {
12,917,638✔
1891
      qError("failed to get table tags from meta, reason:%s, suid:%" PRIu64, tstrerror(code), pListInfo->idInfo.suid);
×
1892
      terrno = code;
×
1893
      QUERY_CHECK_CODE(code, lino, end);
×
1894
    }
1895
  }
1896

1897
  qDebug("final pUidTagList size:%d", (int32_t)taosArrayGetSize(pUidTagList));
13,334,318✔
1898

1899
  int32_t numOfTables = taosArrayGetSize(pUidTagList);
13,336,903✔
1900
  if (numOfTables == 0) {
13,336,708✔
1901
    goto end;
1,894,985✔
1902
  }
1903

1904
  SArray* pColList = NULL;
11,441,723✔
1905
  code = qGetColumnsFromNodeList(pTagCond, false, &pColList); 
11,441,931✔
1906
  if (code != TSDB_CODE_SUCCESS) {
11,439,762✔
1907
    goto end;
×
1908
  }
1909
  pResBlock = createTagValBlockForFilter(pColList, numOfTables, pUidTagList, pVnode, pAPI);
11,439,762✔
1910
  taosArrayDestroy(pColList);
11,438,148✔
1911
  if (pResBlock == NULL) {
11,440,471✔
1912
    code = terrno;
×
1913
    QUERY_CHECK_CODE(code, lino, end);
×
1914
  }
1915

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

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

1923
  void* tmp = taosArrayPush(pBlockList, &pResBlock);
11,441,931✔
1924
  QUERY_CHECK_NULL(tmp, code, lino, end, terrno);
11,441,931✔
1925

1926
  code = createResultData(&type, numOfTables, &output);
11,441,931✔
1927
  if (code != TSDB_CODE_SUCCESS) {
11,438,463✔
1928
    terrno = code;
×
1929
    QUERY_CHECK_CODE(code, lino, end);
×
1930
  }
1931

1932
  gTaskScalarExtra.pStreamInfo = pStreamInfo;
11,438,463✔
1933
  gTaskScalarExtra.pStreamRange = NULL;
11,438,463✔
1934
  code = scalarCalculate(pTagCond, pBlockList, &output, &gTaskScalarExtra);
11,436,146✔
1935
  if (code != TSDB_CODE_SUCCESS) {
11,437,696✔
1936
    qError("failed to calculate scalar, reason:%s", tstrerror(code));
1,100✔
1937
    terrno = code;
1,100✔
1938
    QUERY_CHECK_CODE(code, lino, end);
1,100✔
1939
  }
1940

1941
  code = doSetQualifiedUid(pListInfo, pUidList, pUidTagList, (bool*)output.columnData->pData, addUid);
11,436,596✔
1942
  if (code != TSDB_CODE_SUCCESS) {
11,440,003✔
1943
    terrno = code;
×
1944
    QUERY_CHECK_CODE(code, lino, end);
×
1945
  }
1946
  *listAdded = true;
11,440,003✔
1947

1948
end:
13,336,206✔
1949
  if (code != TSDB_CODE_SUCCESS) {
13,335,886✔
1950
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
1,100✔
1951
  }
1952
  blockDataDestroy(pResBlock);
13,335,886✔
1953
  taosArrayDestroy(pBlockList);
13,336,365✔
1954
  taosArrayDestroyEx(pUidTagList, freeItem);
13,335,592✔
1955

1956
  colDataDestroy(output.columnData);
13,336,512✔
1957
  taosMemoryFreeClear(output.columnData);
13,336,842✔
1958
  return code;
13,336,194✔
1959
}
1960

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

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

1989
  return DEAL_RES_CONTINUE;
26,733✔
1990
}
1991

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

1998
  col_id_t colId = pColNode->colId;
42,834✔
1999
  void* _tmp = taosArrayPush(pColIdArray, &colId);
42,834✔
2000
}
42,834✔
2001

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

2015
  if (nodeType(pTagCond) == QUERY_NODE_OPERATOR) {
21,417✔
2016
    extractTagColId((SOperatorNode*)pTagCond, *pTagColIds);
×
2017
  } else if (nodeType(pTagCond) == QUERY_NODE_LOGIC_CONDITION) {
21,417✔
2018
    SNode* pChild = NULL;
21,417✔
2019
    FOREACH(pChild, ((SLogicConditionNode*)pTagCond)->pParameterList) {
64,251✔
2020
      extractTagColId((SOperatorNode*)pChild, *pTagColIds);
42,834✔
2021
    }
2022
  }
2023

2024
  taosArraySort(*pTagColIds, compareUint16Val);
21,417✔
2025

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

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

2051
static EDealRes canOptimizeTagCondFilter(SNode* pTagCond, void* pContext) {
199,650✔
2052
  if (NULL == pTagCond) {
199,650✔
2053
    *(bool*)pContext = false;
×
2054
    return DEAL_RES_END;
×
2055
  }
2056
  if (nodeType(pTagCond) == QUERY_NODE_VALUE ||
199,650✔
2057
    nodeType(pTagCond) == QUERY_NODE_COLUMN) {
132,495✔
2058
    return DEAL_RES_CONTINUE;
109,989✔
2059
  }
2060
  if (nodeType(pTagCond) == QUERY_NODE_OPERATOR &&
89,661✔
2061
    ((SOperatorNode*)pTagCond)->opType == OP_TYPE_EQUAL) {
43,923✔
2062
    return DEAL_RES_CONTINUE;
42,834✔
2063
  }
2064
  if (nodeType(pTagCond) == QUERY_NODE_LOGIC_CONDITION &&
46,827✔
2065
    ((SLogicConditionNode*)pTagCond)->condType == LOGIC_COND_TYPE_AND) {
21,417✔
2066
    return DEAL_RES_CONTINUE;
21,417✔
2067
  }
2068
  if (nodeType(pTagCond) == QUERY_NODE_FUNCTION &&
49,731✔
2069
    fmIsStreamPesudoColVal(((SFunctionNode*)pTagCond)->funcId)) {
24,321✔
2070
    return DEAL_RES_CONTINUE;
24,321✔
2071
  }
2072
  *(bool*)pContext = false;
1,089✔
2073
  return DEAL_RES_END;
1,089✔
2074
}
2075

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

2083
  pListInfo->idInfo.suid = pScanNode->suid;
214,198,604✔
2084
  pListInfo->idInfo.tableType = pScanNode->tableType;
214,149,575✔
2085

2086
  SArray* pUidList = taosArrayInit(8, sizeof(uint64_t));
214,153,404✔
2087
  QUERY_CHECK_NULL(pUidList, code, lino, _error, terrno);
214,140,901✔
2088

2089
  SIdxFltStatus status = SFLT_NOT_INDEX;
214,140,901✔
2090
  char*   pTagCondKey = NULL;
214,153,106✔
2091
  int32_t tagCondKeyLen;
214,085,051✔
2092
  SArray* pTagColIds = NULL;
214,092,332✔
2093
  char*   pPayload = NULL;
214,116,814✔
2094
  qTrace("getTableList called, suid:%" PRIu64
214,116,814✔
2095
    ", tagCond:%p, tagIndexCond:%p, %d %d", pScanNode->suid, pTagCond,
2096
    pTagIndexCond, pScanNode->tableType, pScanNode->virtualStableScan);
2097
  if (pScanNode->tableType != TSDB_SUPER_TABLE && !pScanNode->virtualStableScan) {
214,116,814✔
2098
    pListInfo->idInfo.uid = pScanNode->uid;
143,168,309✔
2099
    if (pStorageAPI->metaFn.isTableExisted(pVnode, pScanNode->uid)) {
143,140,188✔
2100
      void* tmp = taosArrayPush(pUidList, &pScanNode->uid);
143,192,415✔
2101
      QUERY_CHECK_NULL(tmp, code, lino, _error, terrno);
143,194,294✔
2102
    }
2103
    code = doFilterByTagCond(pListInfo, pUidList, pTagCond, pVnode, status, pStorageAPI, false, &listAdded, pStreamInfo);
143,247,508✔
2104
    QUERY_CHECK_CODE(code, lino, _end);
143,238,772✔
2105
  } else {
2106
    bool      isStream = (pStreamInfo != NULL);
70,894,298✔
2107
    bool      hasTagCond = (pTagCond != NULL);
70,894,298✔
2108
    bool      canCacheTagEqCondFilter = false;
70,894,298✔
2109
    T_MD5_CTX context = {0};
70,958,998✔
2110

2111
    qTrace("start to get table list by tag filter, suid:%" PRIu64
71,023,067✔
2112
      ",tsStableTagFilterCache:%d, tsTagFilterCache:%d", 
2113
      pScanNode->suid, tsStableTagFilterCache, tsTagFilterCache);
2114

2115
    bool acquired = false;
71,023,067✔
2116
    // first, check whether we can use stable tag filter cache
2117
    if (tsStableTagFilterCache && isStream && hasTagCond) {
70,943,589✔
2118
      canCacheTagEqCondFilter = true;
22,506✔
2119
      nodesWalkExpr(pTagCond, canOptimizeTagCondFilter,
22,506✔
2120
        (void*)&canCacheTagEqCondFilter);
2121
    }
2122
    if (canCacheTagEqCondFilter) {
70,854,075✔
2123
      qDebug("%s, stable tag filter condition can be optimized", idstr);
21,417✔
2124
      if (((SStreamRuntimeFuncInfo*)pStreamInfo)->hasPlaceHolder) {
21,417✔
2125
        SNode* tmp = NULL;
21,417✔
2126
        code = nodesCloneNode((SNode*)pTagCond, &tmp);
21,417✔
2127
        QUERY_CHECK_CODE(code, lino, _error);
21,417✔
2128

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

2143
      code = buildTagCondKey(
21,417✔
2144
        pTagCond, &pTagCondKey, &tagCondKeyLen, &pTagColIds);
2145
      QUERY_CHECK_CODE(code, lino, _error);
21,417✔
2146
      code = pStorageAPI->metaFn.getStableCachedTableList(
21,417✔
2147
        pVnode, pScanNode->suid, pTagCondKey, tagCondKeyLen,
21,417✔
2148
        context.digest, tListLen(context.digest), pUidList, &acquired);
2149
      QUERY_CHECK_CODE(code, lino, _error);
21,417✔
2150
    } else if (tsTagFilterCache) {
70,832,658✔
2151
      // second, try to use normal tag filter cache
2152
      qDebug("%s using normal tag filter cache", idstr);
65,811✔
2153
      if (pStreamInfo != NULL && ((SStreamRuntimeFuncInfo*)pStreamInfo)->hasPlaceHolder) {
68,223✔
2154
        SNode* tmp = NULL;
2,412✔
2155
        code = nodesCloneNode((SNode*)pTagCond, &tmp);
2,412✔
2156
        QUERY_CHECK_CODE(code, lino, _error);
2,412✔
2157

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

2193
    if (!pTagCond) {  // no tag filter condition exists, let's fetch all tables of this super table
70,943,964✔
2194
      code = pStorageAPI->metaFn.getChildTableList(pVnode, pScanNode->suid, pUidList);
57,884,012✔
2195
      QUERY_CHECK_CODE(code, lino, _error);
57,886,555✔
2196
      qTrace("no tag filter, get all child tables, numOfTables:%d", (int32_t)taosArrayGetSize(pUidList));
57,886,555✔
2197
    } else {
2198
      // failed to find the result in the cache, let try to calculate the results
2199
      if (pTagIndexCond) {
13,059,952✔
2200
        void* pIndex = pStorageAPI->metaFn.getInvertIndex(pVnode);
4,515,974✔
2201

2202
        SIndexMetaArg metaArg = {.metaEx = pVnode,
4,516,042✔
2203
                                 .idx = pStorageAPI->metaFn.storeGetIndexInfo(pVnode),
4,515,974✔
2204
                                 .ivtIdx = pIndex,
2205
                                 .suid = pScanNode->uid};
4,515,974✔
2206

2207
        status = SFLT_NOT_INDEX;
4,515,974✔
2208
        code = doFilterTag(pTagIndexCond, &metaArg, pUidList, &status, &pStorageAPI->metaFilter);
4,515,974✔
2209
        if (code != 0 || status == SFLT_NOT_INDEX) {  // temporarily disable it for performance sake
4,511,752✔
2210
          qDebug("failed to get tableIds from index, suid:%" PRIu64 ", uidListSize:%d", pScanNode->uid, (int32_t)taosArrayGetSize(pUidList));
1,081,520✔
2211
        } else {
2212
          qDebug("succ to get filter result, table num: %d", (int)taosArrayGetSize(pUidList));
3,430,232✔
2213
        }
2214
      }
2215
    }
2216
    qTrace("after index filter, pTagCond:%p uidListSize:%d", pTagCond, (int32_t)taosArrayGetSize(pUidList));
70,930,074✔
2217
    code = doFilterByTagCond(pListInfo, pUidList, pTagCond, pVnode, status,
70,963,715✔
2218
      pStorageAPI, tsTagFilterCache || tsStableTagFilterCache,
70,963,715✔
2219
      &listAdded, pStreamInfo);
2220
    QUERY_CHECK_CODE(code, lino, _error);
70,954,317✔
2221

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

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

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

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

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

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

2276

2277
      digest[0] = 1;
16,311✔
2278
      memcpy(digest + 1, context.digest, tListLen(context.digest));
16,311✔
2279
    }
2280
  }
2281

2282
_end:
214,249,299✔
2283
  if (!listAdded) {
214,140,284✔
2284
    numOfTables = taosArrayGetSize(pUidList);
202,779,422✔
2285
    for (int i = 0; i < numOfTables; i++) {
618,828,583✔
2286
      void* tmp = taosArrayGet(pUidList, i);
415,972,791✔
2287
      QUERY_CHECK_NULL(tmp, code, lino, _error, terrno);
415,992,982✔
2288
      STableKeyInfo info = {.uid = *(uint64_t*)tmp, .groupId = 0};
415,992,982✔
2289

2290
      void* p = taosArrayPush(pListInfo->pTableList, &info);
415,954,567✔
2291
      if (p == NULL) {
416,056,418✔
2292
        taosArrayDestroy(pUidList);
×
2293
        return terrno;
×
2294
      }
2295

2296
      qTrace("tagfilter get uid:%" PRIu64 ", %s", info.uid, idstr);
416,056,418✔
2297
    }
2298
  }
2299

2300
  qDebug("%s, table list with %d uids built", idstr, (int32_t)numOfTables);
214,216,654✔
2301

2302
_error:
214,229,014✔
2303
  taosArrayDestroy(pUidList);
214,268,997✔
2304
  taosArrayDestroy(pTagColIds);
214,246,921✔
2305
  taosMemFreeClear(pTagCondKey);
214,264,120✔
2306
  taosMemFreeClear(pPayload);
214,264,120✔
2307
  if (code != TSDB_CODE_SUCCESS) {
214,264,120✔
2308
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
1,100✔
2309
  }
2310
  return code;
214,203,466✔
2311
}
2312

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

2435
  return TSDB_CODE_SUCCESS;
×
2436
}
2437

2438
SArray* makeColumnArrayFromList(SNodeList* pNodeList) {
7,566,588✔
2439
  if (!pNodeList) {
7,566,588✔
2440
    return NULL;
×
2441
  }
2442

2443
  size_t  numOfCols = LIST_LENGTH(pNodeList);
7,566,588✔
2444
  SArray* pList = taosArrayInit(numOfCols, sizeof(SColumn));
7,567,041✔
2445
  if (pList == NULL) {
7,566,122✔
2446
    return NULL;
×
2447
  }
2448

2449
  for (int32_t i = 0; i < numOfCols; ++i) {
17,037,531✔
2450
    SColumnNode* pColNode = (SColumnNode*)nodesListGetNode(pNodeList, i);
9,472,023✔
2451
    if (!pColNode) {
9,474,007✔
2452
      taosArrayDestroy(pList);
×
2453
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(TSDB_CODE_QRY_EXECUTOR_INTERNAL_ERROR));
×
2454
      return NULL;
×
2455
    }
2456

2457
    // todo extract method
2458
    SColumn c = {0};
9,474,007✔
2459
    c.slotId = pColNode->slotId;
9,472,356✔
2460
    c.colId = pColNode->colId;
9,472,809✔
2461
    c.type = pColNode->node.resType.type;
9,474,007✔
2462
    c.bytes = pColNode->node.resType.bytes;
9,474,007✔
2463
    c.precision = pColNode->node.resType.precision;
9,473,408✔
2464
    c.scale = pColNode->node.resType.scale;
9,471,890✔
2465

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

2474
  return pList;
7,565,508✔
2475
}
2476

2477
int32_t extractColMatchInfo(SNodeList* pNodeList, SDataBlockDescNode* pOutputNodeList, int32_t* numOfOutputCols,
263,001,570✔
2478
                            int32_t type, SColMatchInfo* pMatchInfo) {
2479
  size_t  numOfCols = LIST_LENGTH(pNodeList);
263,001,570✔
2480
  int32_t code = TSDB_CODE_SUCCESS;
263,032,554✔
2481
  int32_t lino = 0;
263,032,554✔
2482

2483
  pMatchInfo->matchType = type;
263,032,554✔
2484

2485
  SArray* pList = taosArrayInit(numOfCols, sizeof(SColMatchItem));
262,996,745✔
2486
  if (pList == NULL) {
262,958,699✔
2487
    code = terrno;
×
2488
    return code;
×
2489
  }
2490

2491
  for (int32_t i = 0; i < numOfCols; ++i) {
1,207,529,740✔
2492
    STargetNode* pNode = (STargetNode*)nodesListGetNode(pNodeList, i);
944,468,858✔
2493
    QUERY_CHECK_NULL(pNode, code, lino, _end, terrno);
944,523,811✔
2494
    if (nodeType(pNode->pExpr) == QUERY_NODE_COLUMN) {
944,523,811✔
2495
      SColumnNode* pColNode = (SColumnNode*)pNode->pExpr;
940,163,127✔
2496

2497
      SColMatchItem c = {.needOutput = true};
940,173,057✔
2498
      c.colId = pColNode->colId;
940,176,630✔
2499
      c.srcSlotId = pColNode->slotId;
940,158,771✔
2500
      c.dstSlotId = pNode->slotId;
940,160,032✔
2501
      c.isPk = pColNode->isPk;
940,162,860✔
2502
      c.dataType = pColNode->node.resType;
940,158,007✔
2503
      void* tmp = taosArrayPush(pList, &c);
940,170,605✔
2504
      QUERY_CHECK_NULL(tmp, code, lino, _end, terrno);
940,170,605✔
2505
    }
2506
  }
2507

2508
  // set the output flag for each column in SColMatchInfo, according to the
2509
  *numOfOutputCols = 0;
263,060,882✔
2510
  int32_t num = LIST_LENGTH(pOutputNodeList->pSlots);
263,071,213✔
2511
  for (int32_t i = 0; i < num; ++i) {
1,310,658,506✔
2512
    SSlotDescNode* pNode = (SSlotDescNode*)nodesListGetNode(pOutputNodeList->pSlots, i);
1,047,612,448✔
2513
    QUERY_CHECK_NULL(pNode, code, lino, _end, terrno);
1,047,656,705✔
2514

2515
    // todo: add reserve flag check
2516
    // it is a column reserved for the arithmetic expression calculation
2517
    if (pNode->slotId >= numOfCols) {
1,047,656,705✔
2518
      (*numOfOutputCols) += 1;
103,071,205✔
2519
      continue;
103,071,609✔
2520
    }
2521

2522
    SColMatchItem* info = NULL;
944,590,833✔
2523
    for (int32_t j = 0; j < taosArrayGetSize(pList); ++j) {
2,147,483,647✔
2524
      info = taosArrayGet(pList, j);
2,147,483,647✔
2525
      QUERY_CHECK_NULL(info, code, lino, _end, terrno);
2,147,483,647✔
2526
      if (info->dstSlotId == pNode->slotId) {
2,147,483,647✔
2527
        break;
939,303,717✔
2528
      }
2529
    }
2530

2531
    if (pNode->output) {
13,125,623✔
2532
      (*numOfOutputCols) += 1;
935,653,451✔
2533
    } else if (info != NULL) {
8,886,967✔
2534
      // select distinct tbname from stb where tbname='abc';
2535
      info->needOutput = false;
8,887,614✔
2536
    }
2537
  }
2538

2539
  pMatchInfo->pList = pList;
263,046,058✔
2540

2541
_end:
262,974,363✔
2542
  if (code != TSDB_CODE_SUCCESS) {
262,974,363✔
2543
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
2544
  }
2545
  return code;
262,986,262✔
2546
}
2547

2548
static SResSchema createResSchema(int32_t type, int32_t bytes, int32_t slotId, int32_t scale, int32_t precision,
782,834,090✔
2549
                                  const char* name) {
2550
  SResSchema s = {0};
782,834,090✔
2551
  s.scale = scale;
782,895,537✔
2552
  s.type = type;
782,895,537✔
2553
  s.bytes = bytes;
782,895,537✔
2554
  s.slotId = slotId;
782,895,537✔
2555
  s.precision = precision;
782,895,537✔
2556
  tstrncpy(s.name, name, tListLen(s.name));
782,895,537✔
2557

2558
  return s;
782,895,537✔
2559
}
2560

2561
static SColumn* createColumn(int32_t blockId, int32_t slotId, int32_t colId, SDataType* pType, EColumnType colType) {
705,397,961✔
2562
  SColumn* pCol = taosMemoryCalloc(1, sizeof(SColumn));
705,397,961✔
2563
  if (pCol == NULL) {
705,270,533✔
2564
    return NULL;
×
2565
  }
2566

2567
  pCol->slotId = slotId;
705,270,533✔
2568
  pCol->colId = colId;
705,282,834✔
2569
  pCol->bytes = pType->bytes;
705,299,510✔
2570
  pCol->type = pType->type;
705,399,818✔
2571
  pCol->scale = pType->scale;
705,394,666✔
2572
  pCol->precision = pType->precision;
705,366,170✔
2573
  pCol->dataBlockId = blockId;
705,401,997✔
2574
  pCol->colType = colType;
705,392,136✔
2575
  return pCol;
705,319,311✔
2576
}
2577

2578
int32_t createExprFromOneNode(SExprInfo* pExp, SNode* pNode, int16_t slotId) {
787,294,889✔
2579
  int32_t code = TSDB_CODE_SUCCESS;
787,294,889✔
2580
  int32_t lino = 0;
787,294,889✔
2581
  pExp->base.numOfParams = 0;
787,294,889✔
2582
  pExp->base.pParam = NULL;
787,391,357✔
2583
  pExp->pExpr = taosMemoryCalloc(1, sizeof(tExprNode));
787,337,020✔
2584
  QUERY_CHECK_NULL(pExp->pExpr, code, lino, _end, terrno);
787,171,016✔
2585

2586
  pExp->pExpr->_function.num = 1;
787,244,911✔
2587
  pExp->pExpr->_function.functionId = -1;
787,318,429✔
2588

2589
  int32_t type = nodeType(pNode);
787,272,704✔
2590
  // it is a project query, or group by column
2591
  if (type == QUERY_NODE_COLUMN) {
787,363,221✔
2592
    pExp->pExpr->nodeType = QUERY_NODE_COLUMN;
473,606,375✔
2593
    SColumnNode* pColNode = (SColumnNode*)pNode;
473,622,418✔
2594

2595
    pExp->base.pParam = taosMemoryCalloc(1, sizeof(SFunctParam));
473,622,418✔
2596
    QUERY_CHECK_NULL(pExp->base.pParam, code, lino, _end, terrno);
473,545,199✔
2597

2598
    pExp->base.numOfParams = 1;
473,557,445✔
2599

2600
    SDataType* pType = &pColNode->node.resType;
473,541,090✔
2601
    pExp->base.resSchema =
2602
        createResSchema(pType->type, pType->bytes, slotId, pType->scale, pType->precision, pColNode->colName);
473,554,720✔
2603

2604
    pExp->base.pParam[0].pCol =
947,052,608✔
2605
        createColumn(pColNode->dataBlockId, pColNode->slotId, pColNode->colId, pType, pColNode->colType);
947,057,746✔
2606
    QUERY_CHECK_NULL(pExp->base.pParam[0].pCol, code, lino, _end, terrno);
473,549,350✔
2607

2608
    pExp->base.pParam[0].type = FUNC_PARAM_TYPE_COLUMN;
473,469,749✔
2609
  } else if (type == QUERY_NODE_VALUE) {
313,756,846✔
2610
    pExp->pExpr->nodeType = QUERY_NODE_VALUE;
15,485,131✔
2611
    SValueNode* pValNode = (SValueNode*)pNode;
15,487,787✔
2612

2613
    pExp->base.pParam = taosMemoryCalloc(1, sizeof(SFunctParam));
15,487,787✔
2614
    QUERY_CHECK_NULL(pExp->base.pParam, code, lino, _end, terrno);
15,486,798✔
2615

2616
    pExp->base.numOfParams = 1;
15,485,534✔
2617

2618
    SDataType* pType = &pValNode->node.resType;
15,485,746✔
2619
    pExp->base.resSchema =
2620
        createResSchema(pType->type, pType->bytes, slotId, pType->scale, pType->precision, pValNode->node.aliasName);
15,487,417✔
2621
    pExp->base.pParam[0].type = FUNC_PARAM_TYPE_VALUE;
15,484,980✔
2622
    code = nodesValueNodeToVariant(pValNode, &pExp->base.pParam[0].param);
15,485,487✔
2623
    QUERY_CHECK_CODE(code, lino, _end);
15,485,342✔
2624
  } else if (type == QUERY_NODE_REMOTE_VALUE) {
298,271,715✔
2625
    SRemoteValueNode* pRemote = (SRemoteValueNode*)pNode;
25,878,123✔
2626
    code = qFetchRemoteValue(gTaskScalarExtra.pSubJobCtx, pRemote->subQIdx, pRemote);
25,878,123✔
2627
    QUERY_CHECK_CODE(code, lino, _end);
25,899,156✔
2628

2629
    pExp->pExpr->nodeType = QUERY_NODE_VALUE;
21,330,417✔
2630
    SValueNode* pValNode = (SValueNode*)pNode;
21,330,859✔
2631

2632
    pExp->base.pParam = taosMemoryCalloc(1, sizeof(SFunctParam));
21,330,859✔
2633
    QUERY_CHECK_NULL(pExp->base.pParam, code, lino, _end, terrno);
21,330,859✔
2634

2635
    pExp->base.numOfParams = 1;
21,330,859✔
2636

2637
    SDataType* pType = &pValNode->node.resType;
21,330,859✔
2638
    pExp->base.resSchema =
2639
        createResSchema(pType->type, pType->bytes, slotId, pType->scale, pType->precision, pValNode->node.aliasName);
21,330,859✔
2640
    pExp->base.pParam[0].type = FUNC_PARAM_TYPE_VALUE;
21,330,859✔
2641
    code = nodesValueNodeToVariant(pValNode, &pExp->base.pParam[0].param);
21,330,859✔
2642
    QUERY_CHECK_CODE(code, lino, _end);
21,330,859✔
2643
  } else if (type == QUERY_NODE_FUNCTION) {
272,393,592✔
2644
    pExp->pExpr->nodeType = QUERY_NODE_FUNCTION;
241,514,106✔
2645
    SFunctionNode* pFuncNode = (SFunctionNode*)pNode;
241,528,249✔
2646

2647
    SDataType* pType = &pFuncNode->node.resType;
241,528,249✔
2648
    pExp->base.resSchema =
2649
        createResSchema(pType->type, pType->bytes, slotId, pType->scale, pType->precision, pFuncNode->node.aliasName);
241,538,324✔
2650
    tExprNode* pExprNode = pExp->pExpr;
241,508,930✔
2651

2652
    pExprNode->_function.functionId = pFuncNode->funcId;
241,495,995✔
2653
    pExprNode->_function.pFunctNode = pFuncNode;
241,539,201✔
2654
    pExprNode->_function.functionType = pFuncNode->funcType;
241,541,507✔
2655

2656
    tstrncpy(pExprNode->_function.functionName, pFuncNode->functionName, tListLen(pExprNode->_function.functionName));
241,512,948✔
2657

2658
    pExp->base.pParamList = pFuncNode->pParameterList;
241,551,804✔
2659
#if 1
2660
    // todo refactor: add the parameter for tbname function
2661
    const char* name = "tbname";
241,541,054✔
2662
    int32_t     len = strlen(name);
241,541,054✔
2663

2664
    if (!pFuncNode->pParameterList && (memcmp(pExprNode->_function.functionName, name, len) == 0) &&
241,541,054✔
2665
        pExprNode->_function.functionName[len] == 0) {
8,967,284✔
2666
      pFuncNode->pParameterList = NULL;
8,966,573✔
2667
      int32_t     code = nodesMakeList(&pFuncNode->pParameterList);
8,964,974✔
2668
      SValueNode* res = NULL;
8,970,817✔
2669
      if (TSDB_CODE_SUCCESS == code) {
8,971,277✔
2670
        code = nodesMakeNode(QUERY_NODE_VALUE, (SNode**)&res);
8,970,770✔
2671
      }
2672
      QUERY_CHECK_CODE(code, lino, _end);
8,971,629✔
2673
      res->node.resType = (SDataType){.bytes = sizeof(int64_t), .type = TSDB_DATA_TYPE_BIGINT};
8,971,629✔
2674
      code = nodesListAppend(pFuncNode->pParameterList, (SNode*)res);
8,969,044✔
2675
      if (code != TSDB_CODE_SUCCESS) {
8,970,205✔
2676
        nodesDestroyNode((SNode*)res);
×
2677
        res = NULL;
×
2678
      }
2679
      QUERY_CHECK_CODE(code, lino, _end);
8,970,205✔
2680
    }
2681
#endif
2682

2683
    int32_t numOfParam = LIST_LENGTH(pFuncNode->pParameterList);
241,530,069✔
2684

2685
    pExp->base.pParam = taosMemoryCalloc(numOfParam, sizeof(SFunctParam));
241,530,366✔
2686
    QUERY_CHECK_NULL(pExp->base.pParam, code, lino, _end, terrno);
241,519,064✔
2687
    pExp->base.numOfParams = numOfParam;
241,510,911✔
2688

2689
    for (int32_t j = 0; j < numOfParam && TSDB_CODE_SUCCESS == code; ++j) {
601,781,060✔
2690
      SNode* p1 = nodesListGetNode(pFuncNode->pParameterList, j);
360,490,281✔
2691
      QUERY_CHECK_NULL(p1, code, lino, _end, terrno);
360,478,430✔
2692
      if (p1->type == QUERY_NODE_COLUMN) {
360,478,430✔
2693
        SColumnNode* pcn = (SColumnNode*)p1;
231,848,483✔
2694

2695
        pExp->base.pParam[j].type = FUNC_PARAM_TYPE_COLUMN;
231,848,483✔
2696
        pExp->base.pParam[j].pCol =
463,647,333✔
2697
            createColumn(pcn->dataBlockId, pcn->slotId, pcn->colId, &pcn->node.resType, pcn->colType);
463,666,416✔
2698
        QUERY_CHECK_NULL(pExp->base.pParam[j].pCol, code, lino, _end, terrno);
231,819,756✔
2699
      } else if (p1->type == QUERY_NODE_VALUE) {
128,631,810✔
2700
        SValueNode* pvn = (SValueNode*)p1;
64,721,111✔
2701
        pExp->base.pParam[j].type = FUNC_PARAM_TYPE_VALUE;
64,721,111✔
2702
        code = nodesValueNodeToVariant(pvn, &pExp->base.pParam[j].param);
64,709,750✔
2703
        QUERY_CHECK_CODE(code, lino, _end);
64,715,422✔
2704
      } else if (p1->type == QUERY_NODE_REMOTE_VALUE) {
63,932,377✔
2705
        SRemoteValueNode* pRemote = (SRemoteValueNode*)p1;
1,313,719✔
2706
        code = qFetchRemoteValue(gTaskScalarExtra.pSubJobCtx, pRemote->subQIdx, pRemote);
1,313,719✔
2707
        QUERY_CHECK_CODE(code, lino, _end);
1,313,719✔
2708

2709
        SValueNode* pvn = (SValueNode*)pRemote;
1,102,554✔
2710
        pExp->base.pParam[j].type = FUNC_PARAM_TYPE_VALUE;
1,102,554✔
2711
        code = nodesValueNodeToVariant(pvn, &pExp->base.pParam[j].param);
1,102,554✔
2712
        QUERY_CHECK_CODE(code, lino, _end);
1,071,381✔
2713
      }
2714
    }
2715
    pExp->pExpr->_function.bindExprID = ((SExprNode*)pNode)->bindExprID;
241,290,779✔
2716
  } else if (type == QUERY_NODE_OPERATOR) {
30,879,486✔
2717
    pExp->pExpr->nodeType = QUERY_NODE_OPERATOR;
25,610,906✔
2718
    SOperatorNode* pOpNode = (SOperatorNode*)pNode;
25,611,661✔
2719

2720
    pExp->base.pParam = taosMemoryCalloc(1, sizeof(SFunctParam));
25,611,661✔
2721
    QUERY_CHECK_NULL(pExp->base.pParam, code, lino, _end, terrno);
25,607,385✔
2722
    pExp->base.numOfParams = 1;
25,605,536✔
2723

2724
    SDataType* pType = &pOpNode->node.resType;
25,605,903✔
2725
    pExp->base.resSchema =
2726
        createResSchema(pType->type, pType->bytes, slotId, pType->scale, pType->precision, pOpNode->node.aliasName);
25,607,873✔
2727
    pExp->pExpr->_optrRoot.pRootNode = pNode;
25,606,927✔
2728
  } else if (type == QUERY_NODE_CASE_WHEN) {
5,269,402✔
2729
    pExp->pExpr->nodeType = QUERY_NODE_OPERATOR;
5,269,881✔
2730
    SCaseWhenNode* pCaseNode = (SCaseWhenNode*)pNode;
5,269,881✔
2731

2732
    pExp->base.pParam = taosMemoryCalloc(1, sizeof(SFunctParam));
5,269,881✔
2733
    QUERY_CHECK_NULL(pExp->base.pParam, code, lino, _end, terrno);
5,269,881✔
2734
    pExp->base.numOfParams = 1;
5,269,881✔
2735

2736
    SDataType* pType = &pCaseNode->node.resType;
5,269,881✔
2737
    pExp->base.resSchema =
2738
        createResSchema(pType->type, pType->bytes, slotId, pType->scale, pType->precision, pCaseNode->node.aliasName);
5,269,881✔
2739
    pExp->pExpr->_optrRoot.pRootNode = pNode;
5,269,881✔
2740
  } else if (type == QUERY_NODE_LOGIC_CONDITION) {
2,319✔
2741
    pExp->pExpr->nodeType = QUERY_NODE_OPERATOR;
1,148✔
2742
    SLogicConditionNode* pCond = (SLogicConditionNode*)pNode;
1,148✔
2743
    pExp->base.pParam = taosMemoryCalloc(1, sizeof(SFunctParam));
1,148✔
2744
    QUERY_CHECK_NULL(pExp->base.pParam, code, lino, _end, terrno);
1,148✔
2745
    pExp->base.numOfParams = 1;
1,148✔
2746
    SDataType* pType = &pCond->node.resType;
1,148✔
2747
    pExp->base.resSchema =
2748
        createResSchema(pType->type, pType->bytes, slotId, pType->scale, pType->precision, pCond->node.aliasName);
1,148✔
2749
    pExp->pExpr->_optrRoot.pRootNode = pNode;
1,148✔
2750
  } else {
2751
    code = TSDB_CODE_QRY_EXECUTOR_INTERNAL_ERROR;
1,375✔
2752
    QUERY_CHECK_CODE(code, lino, _end);
1,375✔
2753
  }
2754
  pExp->pExpr->relatedTo = ((SExprNode*)pNode)->relatedTo;
782,526,695✔
2755
_end:
787,259,729✔
2756
  if (code != TSDB_CODE_SUCCESS) {
787,259,729✔
2757
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
4,779,904✔
2758
  }
2759
  return code;
787,291,010✔
2760
}
2761

2762
int32_t createExprFromTargetNode(SExprInfo* pExp, STargetNode* pTargetNode) {
787,263,586✔
2763
  return createExprFromOneNode(pExp, pTargetNode->pExpr, pTargetNode->slotId);
787,263,586✔
2764
}
2765

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

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

2784
  return pExprs;
×
2785
}
2786

2787
int32_t createExprInfo(SNodeList* pNodeList, SNodeList* pGroupKeys, SExprInfo** pExprInfo, int32_t* numOfExprs) {
349,800,934✔
2788
  QRY_PARAM_CHECK(pExprInfo);
349,800,934✔
2789

2790
  int32_t code = 0;
349,847,048✔
2791
  int32_t numOfFuncs = LIST_LENGTH(pNodeList);
349,847,048✔
2792
  int32_t numOfGroupKeys = 0;
349,835,979✔
2793
  if (pGroupKeys != NULL) {
349,835,979✔
2794
    numOfGroupKeys = LIST_LENGTH(pGroupKeys);
33,767,118✔
2795
  }
2796

2797
  *numOfExprs = numOfFuncs + numOfGroupKeys;
349,836,634✔
2798
  if (*numOfExprs == 0) {
349,777,203✔
2799
    return code;
41,713,840✔
2800
  }
2801

2802
  SExprInfo* pExprs = taosMemoryCalloc(*numOfExprs, sizeof(SExprInfo));
308,143,709✔
2803
  if (pExprs == NULL) {
307,972,441✔
2804
    return terrno;
×
2805
  }
2806

2807
  for (int32_t i = 0; i < (*numOfExprs); ++i) {
1,090,289,352✔
2808
    STargetNode* pTargetNode = NULL;
787,199,723✔
2809
    if (i < numOfFuncs) {
787,199,723✔
2810
      pTargetNode = (STargetNode*)nodesListGetNode(pNodeList, i);
742,358,617✔
2811
    } else {
2812
      pTargetNode = (STargetNode*)nodesListGetNode(pGroupKeys, i - numOfFuncs);
44,841,106✔
2813
    }
2814
    if (!pTargetNode) {
787,283,003✔
2815
      destroyExprInfo(pExprs, *numOfExprs);
×
2816
      taosMemoryFreeClear(pExprs);
×
2817
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
2818
      return terrno;
×
2819
    }
2820

2821
    SExprInfo* pExp = &pExprs[i];
787,283,003✔
2822
    code = createExprFromTargetNode(pExp, pTargetNode);
787,243,507✔
2823
    if (code != TSDB_CODE_SUCCESS) {
787,096,815✔
2824
      destroyExprInfo(pExprs, *numOfExprs);
4,779,904✔
2825
      taosMemoryFreeClear(pExprs);
4,779,904✔
2826
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(code));
4,779,904✔
2827
      return code;
4,779,904✔
2828
    }
2829
  }
2830

2831
  *pExprInfo = pExprs;
303,266,838✔
2832
  return code;
303,309,433✔
2833
}
2834

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

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

2848
  SArray* pValCtxArray = NULL;
325,232,409✔
2849
  for (int32_t i = numOfOutput - 1; i > 0; --i) {  // select Func is at the end of the list
799,612,638✔
2850
    int32_t funcIdx = pCtx[i].pExpr->pExpr->_function.bindExprID;
474,394,777✔
2851
    if (funcIdx > 0) {
474,399,760✔
2852
      if (pValCtxArray == NULL) {
1,850,600✔
2853
        // the end of the list is the select function of biggest index
2854
        pValCtxArray = taosArrayInit_s(sizeof(SSubsidiaryResInfo*), funcIdx);
1,329,193✔
2855
        if (pValCtxArray == NULL) {
1,326,658✔
2856
          return terrno;
×
2857
        }
2858
      }
2859
      if (funcIdx > pValCtxArray->size) {
1,848,065✔
2860
        qError("funcIdx:%d is out of range", funcIdx);
×
2861
        taosArrayDestroyP(pValCtxArray, deleteSubsidiareCtx);
×
2862
        return TSDB_CODE_QRY_EXECUTOR_INTERNAL_ERROR;
×
2863
      }
2864
      SSubsidiaryResInfo* pSubsidiary = &pCtx[i].subsidiaries;
1,851,107✔
2865
      pSubsidiary->pCtx = taosMemoryCalloc(numOfOutput, POINTER_BYTES);
1,849,586✔
2866
      if (pSubsidiary->pCtx == NULL) {
1,851,614✔
2867
        taosArrayDestroyP(pValCtxArray, deleteSubsidiareCtx);
×
2868
        return terrno;
×
2869
      }
2870
      pSubsidiary->num = 0;
1,852,121✔
2871
      taosArraySet(pValCtxArray, funcIdx - 1, &pSubsidiary);
1,852,121✔
2872
    }
2873
  }
2874

2875
  SqlFunctionCtx*  p = NULL;
325,217,861✔
2876
  SqlFunctionCtx** pValCtx = NULL;
325,217,861✔
2877
  if (pValCtxArray == NULL) {
325,217,861✔
2878
    pValCtx = taosMemoryCalloc(numOfOutput, POINTER_BYTES);
324,056,429✔
2879
    if (pValCtx == NULL) {
324,216,070✔
2880
      QUERY_CHECK_CODE(terrno, lino, _end);
×
2881
    }
2882
  }
2883

2884
  for (int32_t i = 0; i < numOfOutput; ++i) {
1,094,138,218✔
2885
    const char* pName = pCtx[i].pExpr->pExpr->_function.functionName;
768,785,409✔
2886
    if ((strcmp(pName, "_select_value") == 0)) {
768,804,869✔
2887
      if (pValCtxArray == NULL) {
5,911,352✔
2888
        pValCtx[num++] = &pCtx[i];
3,308,684✔
2889
      } else {
2890
        int32_t bindFuncIndex = pCtx[i].pExpr->pExpr->relatedTo;  // start from index 1;
2,602,668✔
2891
        if (bindFuncIndex > 0) {                                  // 0 is default index related to the select function
2,602,668✔
2892
          bindFuncIndex -= 1;
2,542,842✔
2893
        }
2894
        SSubsidiaryResInfo** pSubsidiary = taosArrayGet(pValCtxArray, bindFuncIndex);
2,602,668✔
2895
        if (pSubsidiary == NULL) {
2,602,668✔
2896
          QUERY_CHECK_CODE(TSDB_CODE_QRY_EXECUTOR_INTERNAL_ERROR, lino, _end);
×
2897
        }
2898
        (*pSubsidiary)->pCtx[(*pSubsidiary)->num] = &pCtx[i];
2,602,668✔
2899
        (*pSubsidiary)->num++;
2,603,175✔
2900
      }
2901
    } else if (fmIsSelectFunc(pCtx[i].functionId)) {
762,893,517✔
2902
      if (pValCtxArray == NULL) {
56,809,185✔
2903
        p = &pCtx[i];
54,552,031✔
2904
      }
2905
    }
2906
  }
2907

2908
  if (p != NULL) {
325,352,809✔
2909
    p->subsidiaries.pCtx = pValCtx;
24,600,296✔
2910
    p->subsidiaries.num = num;
24,594,352✔
2911
  } else {
2912
    taosMemoryFreeClear(pValCtx);
300,752,513✔
2913
  }
2914

2915
_end:
1,354,562✔
2916
  if (code != TSDB_CODE_SUCCESS) {
325,353,895✔
2917
    taosArrayDestroyP(pValCtxArray, deleteSubsidiareCtx);
×
2918
    taosMemoryFreeClear(pValCtx);
×
2919
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
2920
  } else {
2921
    taosArrayDestroy(pValCtxArray);
325,353,895✔
2922
  }
2923
  return code;
325,334,780✔
2924
}
2925

2926
SqlFunctionCtx* createSqlFunctionCtx(SExprInfo* pExprInfo, int32_t numOfOutput, int32_t** rowEntryInfoOffset,
325,385,657✔
2927
                                     SFunctionStateStore* pStore) {
2928
  int32_t         code = TSDB_CODE_SUCCESS;
325,385,657✔
2929
  int32_t         lino = 0;
325,385,657✔
2930
  SqlFunctionCtx* pFuncCtx = (SqlFunctionCtx*)taosMemoryCalloc(numOfOutput, sizeof(SqlFunctionCtx));
325,385,657✔
2931
  if (pFuncCtx == NULL) {
325,304,481✔
2932
    return NULL;
×
2933
  }
2934

2935
  *rowEntryInfoOffset = taosMemoryCalloc(numOfOutput, sizeof(int32_t));
325,304,481✔
2936
  if (*rowEntryInfoOffset == 0) {
325,412,760✔
2937
    taosMemoryFreeClear(pFuncCtx);
×
2938
    return NULL;
×
2939
  }
2940

2941
  for (int32_t i = 0; i < numOfOutput; ++i) {
1,094,251,372✔
2942
    SExprInfo* pExpr = &pExprInfo[i];
768,892,690✔
2943

2944
    SExprBasicInfo* pFunct = &pExpr->base;
768,824,189✔
2945
    SqlFunctionCtx* pCtx = &pFuncCtx[i];
768,863,339✔
2946

2947
    pCtx->functionId = -1;
768,891,123✔
2948
    pCtx->pExpr = pExpr;
768,877,554✔
2949

2950
    if (pExpr->pExpr->nodeType == QUERY_NODE_FUNCTION) {
768,887,852✔
2951
      SFuncExecEnv env = {0};
239,947,748✔
2952
      pCtx->functionId = pExpr->pExpr->_function.pFunctNode->funcId;
239,948,013✔
2953
      pCtx->isPseudoFunc = fmIsWindowPseudoColumnFunc(pCtx->functionId) || fmIsPlaceHolderFunc(pCtx->functionId);
239,935,273✔
2954
      pCtx->isNotNullFunc = fmIsNotNullOutputFunc(pCtx->functionId);
239,929,117✔
2955

2956
      bool isUdaf = fmIsUserDefinedFunc(pCtx->functionId);
239,923,682✔
2957
      if (fmIsAggFunc(pCtx->functionId) || fmIsIndefiniteRowsFunc(pCtx->functionId)) {
394,510,672✔
2958
        if (!isUdaf) {
154,624,181✔
2959
          code = fmGetFuncExecFuncs(pCtx->functionId, &pCtx->fpSet);
154,582,199✔
2960
          QUERY_CHECK_CODE(code, lino, _end);
154,558,814✔
2961
        } else {
2962
          char* udfName = pExpr->pExpr->_function.pFunctNode->functionName;
41,982✔
2963
          pCtx->udfName = taosStrdup(udfName);
41,982✔
2964
          QUERY_CHECK_NULL(pCtx->udfName, code, lino, _end, terrno);
41,982✔
2965

2966
          code = fmGetUdafExecFuncs(pCtx->functionId, &pCtx->fpSet);
41,982✔
2967
          QUERY_CHECK_CODE(code, lino, _end);
41,982✔
2968
        }
2969
        bool tmp = pCtx->fpSet.getEnv(pExpr->pExpr->_function.pFunctNode, &env);
154,600,796✔
2970
        if (!tmp) {
154,605,036✔
2971
          code = terrno;
×
2972
          QUERY_CHECK_CODE(code, lino, _end);
×
2973
        }
2974
      } else {
2975
        if (fmIsPlaceHolderFunc(pCtx->functionId)) {
85,295,996✔
2976
          code = fmGetStreamPesudoFuncEnv(pCtx->functionId, pExpr->base.pParamList, &env);
7,836,738✔
2977
          QUERY_CHECK_CODE(code, lino, _end);
7,837,178✔
2978
        }      
2979
        
2980
        code = fmGetScalarFuncExecFuncs(pCtx->functionId, &pCtx->sfp);
85,299,961✔
2981
        if (code != TSDB_CODE_SUCCESS && isUdaf) {
85,303,617✔
2982
          code = TSDB_CODE_SUCCESS;
26,280✔
2983
        }
2984
        QUERY_CHECK_CODE(code, lino, _end);
85,303,617✔
2985

2986
        if (pCtx->sfp.getEnv != NULL) {
85,303,617✔
2987
          bool tmp = pCtx->sfp.getEnv(pExpr->pExpr->_function.pFunctNode, &env);
16,555,405✔
2988
          if (!tmp) {
16,556,404✔
2989
            code = terrno;
×
2990
            QUERY_CHECK_CODE(code, lino, _end);
×
2991
          }
2992
        }
2993
      }
2994
      pCtx->resDataInfo.interBufSize = env.calcMemSize;
239,911,920✔
2995
    } else if (pExpr->pExpr->nodeType == QUERY_NODE_COLUMN || pExpr->pExpr->nodeType == QUERY_NODE_OPERATOR ||
528,887,566✔
2996
               pExpr->pExpr->nodeType == QUERY_NODE_VALUE) {
36,769,076✔
2997
      // for simple column, the result buffer needs to hold at least one element.
2998
      pCtx->resDataInfo.interBufSize = pFunct->resSchema.bytes;
528,980,024✔
2999
    }
3000

3001
    pCtx->input.numOfInputCols = pFunct->numOfParams;
768,946,623✔
3002
    pCtx->input.pData = taosMemoryCalloc(pFunct->numOfParams, POINTER_BYTES);
768,898,510✔
3003
    QUERY_CHECK_NULL(pCtx->input.pData, code, lino, _end, terrno);
768,912,168✔
3004
    pCtx->input.pColumnDataAgg = taosMemoryCalloc(pFunct->numOfParams, POINTER_BYTES);
768,900,691✔
3005
    QUERY_CHECK_NULL(pCtx->input.pColumnDataAgg, code, lino, _end, terrno);
768,899,676✔
3006

3007
    pCtx->pTsOutput = NULL;
768,889,619✔
3008
    pCtx->resDataInfo.bytes = pFunct->resSchema.bytes;
768,881,746✔
3009
    pCtx->resDataInfo.type = pFunct->resSchema.type;
768,905,809✔
3010
    pCtx->order = TSDB_ORDER_ASC;
768,824,456✔
3011
    pCtx->start.key = INT64_MIN;
768,913,835✔
3012
    pCtx->end.key = INT64_MIN;
768,882,096✔
3013
    pCtx->numOfParams = pExpr->base.numOfParams;
768,893,770✔
3014
    pCtx->param = pFunct->pParam;
768,917,082✔
3015
    pCtx->saveHandle.currentPage = -1;
768,881,672✔
3016
    pCtx->pStore = pStore;
768,823,635✔
3017
    pCtx->hasWindowOrGroup = false;
768,881,456✔
3018
    pCtx->needCleanup = false;
768,877,678✔
3019
  }
3020

3021
  for (int32_t i = 1; i < numOfOutput; ++i) {
799,718,714✔
3022
    (*rowEntryInfoOffset)[i] = (int32_t)((*rowEntryInfoOffset)[i - 1] + sizeof(SResultRowEntryInfo) +
948,879,601✔
3023
                                         pFuncCtx[i - 1].resDataInfo.interBufSize);
474,444,841✔
3024
  }
3025

3026
  code = setSelectValueColumnInfo(pFuncCtx, numOfOutput);
325,273,776✔
3027
  QUERY_CHECK_CODE(code, lino, _end);
325,344,448✔
3028

3029
_end:
325,344,448✔
3030
  if (code != TSDB_CODE_SUCCESS) {
325,298,744✔
3031
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
3032
    for (int32_t i = 0; i < numOfOutput; ++i) {
×
3033
      taosMemoryFree(pFuncCtx[i].input.pData);
×
3034
      taosMemoryFree(pFuncCtx[i].input.pColumnDataAgg);
×
3035
    }
3036
    taosMemoryFreeClear(*rowEntryInfoOffset);
×
3037
    taosMemoryFreeClear(pFuncCtx);
×
3038

3039
    terrno = code;
×
3040
    return NULL;
×
3041
  }
3042
  return pFuncCtx;
325,298,744✔
3043
}
3044

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

3051
  int32_t i = 0, j = 0;
9,830,835✔
3052
  while (i < numOfSrcCols && j < taosArrayGetSize(pColMatchInfo)) {
90,979,486✔
3053
    SColumnInfoData* p = taosArrayGet(pCols, i);
81,148,391✔
3054
    if (!p) {
81,148,011✔
3055
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
3056
      return terrno;
×
3057
    }
3058
    SColMatchItem* pmInfo = taosArrayGet(pColMatchInfo, j);
81,148,011✔
3059
    if (!pmInfo) {
81,148,542✔
3060
      return terrno;
×
3061
    }
3062

3063
    if (p->info.colId == pmInfo->colId) {
81,148,542✔
3064
      SColumnInfoData* pDst = taosArrayGet(pBlock->pDataBlock, pmInfo->dstSlotId);
72,428,930✔
3065
      if (!pDst) {
72,428,237✔
3066
        return terrno;
×
3067
      }
3068
      code = colDataAssign(pDst, p, pBlock->info.rows, &pBlock->info);
72,428,237✔
3069
      if (code != TSDB_CODE_SUCCESS) {
72,429,039✔
3070
        qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(code));
×
3071
        return code;
×
3072
      }
3073
      i++;
72,429,039✔
3074
      j++;
72,429,039✔
3075
    } else if (p->info.colId < pmInfo->colId) {
8,719,612✔
3076
      i++;
8,719,612✔
3077
    } else {
3078
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(TSDB_CODE_QRY_EXECUTOR_INTERNAL_ERROR));
×
3079
      return TSDB_CODE_QRY_EXECUTOR_INTERNAL_ERROR;
×
3080
    }
3081
  }
3082
  return code;
9,830,456✔
3083
}
3084

3085
SInterval extractIntervalInfo(const STableScanPhysiNode* pTableScanNode) {
188,697,727✔
3086
  SInterval interval = {
377,281,043✔
3087
      .interval = pTableScanNode->interval,
188,551,334✔
3088
      .sliding = pTableScanNode->sliding,
188,567,705✔
3089
      .intervalUnit = pTableScanNode->intervalUnit,
188,719,660✔
3090
      .slidingUnit = pTableScanNode->slidingUnit,
188,637,995✔
3091
      .offset = pTableScanNode->offset,
188,628,951✔
3092
      .precision = pTableScanNode->scan.node.pOutputDataBlockDesc->precision,
188,596,703✔
3093
      .timeRange = pTableScanNode->scanRange,
3094
  };
3095
  calcIntervalAutoOffset(&interval);
188,594,310✔
3096

3097
  return interval;
188,558,801✔
3098
}
3099

3100
SColumn extractColumnFromColumnNode(SColumnNode* pColNode) {
49,808,123✔
3101
  SColumn c = {0};
49,808,123✔
3102

3103
  c.slotId = pColNode->slotId;
49,808,123✔
3104
  c.colId = pColNode->colId;
49,807,663✔
3105
  c.type = pColNode->node.resType.type;
49,808,777✔
3106
  c.bytes = pColNode->node.resType.bytes;
49,805,560✔
3107
  c.scale = pColNode->node.resType.scale;
49,806,180✔
3108
  c.precision = pColNode->node.resType.precision;
49,803,838✔
3109
  return c;
49,803,095✔
3110
}
3111

3112

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

3127
  if (cond->skey > cond->ekey || range->skey > range->ekey) {
1,629,910✔
3128
    *twindow = extTwindows[0] = extTwindows[1] = TSWINDOW_DESC_INITIALIZER;
2,990✔
3129
    return code;
2,990✔
3130
  }
3131

3132
  if (range->ekey < cond->skey) {
1,626,920✔
3133
    extTwindows[1] = *cond;
247,667✔
3134
    *twindow = extTwindows[0] = TSWINDOW_DESC_INITIALIZER;
247,667✔
3135
    return code;
247,667✔
3136
  }
3137

3138
  if (cond->ekey < range->skey) {
1,379,253✔
3139
    extTwindows[0] = *cond;
177,920✔
3140
    *twindow = extTwindows[1] = TSWINDOW_DESC_INITIALIZER;
177,920✔
3141
    return code;
177,920✔
3142
  }
3143

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

3151
  return code;
1,201,333✔
3152
}
3153

3154
int32_t initQueryTableDataCond(SQueryTableDataCond* pCond, const STableScanPhysiNode* pTableScanNode,
210,865,363✔
3155
                               const SReadHandle* readHandle, bool applyExtWin) {
3156
  int32_t code = 0;                             
210,865,363✔
3157
  pCond->order = pTableScanNode->scanSeq[0] > 0 ? TSDB_ORDER_ASC : TSDB_ORDER_DESC;
210,865,363✔
3158
  pCond->numOfCols = LIST_LENGTH(pTableScanNode->scan.pScanCols);
210,889,519✔
3159

3160
  pCond->colList = taosMemoryCalloc(pCond->numOfCols, sizeof(SColumnInfo));
210,835,420✔
3161
  if (!pCond->colList) {
210,816,596✔
3162
    return terrno;
×
3163
  }
3164
  pCond->pSlotList = taosMemoryMalloc(sizeof(int32_t) * pCond->numOfCols);
210,818,357✔
3165
  if (pCond->pSlotList == NULL) {
210,821,617✔
3166
    taosMemoryFreeClear(pCond->colList);
×
3167
    return terrno;
×
3168
  }
3169

3170
  // TODO: get it from stable scan node
3171
  pCond->twindows = pTableScanNode->scanRange;
210,782,424✔
3172
  pCond->suid = pTableScanNode->scan.suid;
210,857,437✔
3173
  pCond->type = TIMEWINDOW_RANGE_CONTAINED;
210,812,262✔
3174
  pCond->startVersion = -1;
210,812,542✔
3175
  pCond->endVersion = -1;
210,776,732✔
3176
  pCond->skipRollup = readHandle->skipRollup;
210,688,889✔
3177
  if (readHandle->winRangeValid) {
210,837,030✔
3178
    pCond->twindows = readHandle->winRange;
334,649✔
3179
  }
3180
  pCond->cacheSttStatis = readHandle->cacheSttStatis;
210,754,194✔
3181
  // allowed read stt file optimization mode
3182
  pCond->notLoadData = (pTableScanNode->dataRequired == FUNC_DATA_REQUIRED_NOT_LOAD) &&
421,592,741✔
3183
                       (pTableScanNode->scan.node.pConditions == NULL) && (pTableScanNode->interval == 0);
210,689,852✔
3184

3185
  int32_t j = 0;
210,794,135✔
3186
  for (int32_t i = 0; i < pCond->numOfCols; ++i) {
993,787,796✔
3187
    STargetNode* pNode = (STargetNode*)nodesListGetNode(pTableScanNode->scan.pScanCols, i);
782,955,621✔
3188
    if (!pNode) {
782,651,255✔
3189
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
3190
      return terrno;
×
3191
    }
3192
    SColumnNode* pColNode = (SColumnNode*)pNode->pExpr;
782,651,255✔
3193
    if (pColNode->colType == COLUMN_TYPE_TAG) {
782,715,337✔
3194
      continue;
×
3195
    }
3196

3197
    pCond->colList[j].type = pColNode->node.resType.type;
782,823,885✔
3198
    pCond->colList[j].bytes = pColNode->node.resType.bytes;
782,847,797✔
3199
    pCond->colList[j].colId = pColNode->colId;
782,892,102✔
3200
    pCond->colList[j].pk = pColNode->isPk;
782,864,997✔
3201

3202
    pCond->pSlotList[j] = pNode->slotId;
782,983,581✔
3203
    j += 1;
782,993,661✔
3204
  }
3205

3206
  pCond->numOfCols = j;
210,901,657✔
3207

3208
  if (applyExtWin) {
210,910,177✔
3209
    if (NULL != pTableScanNode->pExtScanRange) {
189,075,620✔
3210
      pCond->type = TIMEWINDOW_RANGE_EXTERNAL;
1,629,910✔
3211
      code = getQueryExtWindow(&pCond->twindows, pTableScanNode->pExtScanRange, &pCond->twindows, pCond->extTwindows);
1,629,910✔
3212
    } else if (readHandle->extWinRangeValid) {
187,285,433✔
3213
      pCond->type = TIMEWINDOW_RANGE_EXTERNAL;
×
3214
      code = getQueryExtWindow(&pCond->twindows, &readHandle->extWinRange, &pCond->twindows, pCond->extTwindows);
×
3215
    }
3216
  }
3217
  
3218
  return code;
210,838,752✔
3219
}
3220

3221
int32_t initQueryTableDataCondWithColArray(SQueryTableDataCond* pCond, SQueryTableDataCond* pOrgCond,
6,703,004✔
3222
                                           const SReadHandle* readHandle, SArray* colArray) {
3223
  int32_t code = TSDB_CODE_SUCCESS;
6,703,004✔
3224
  int32_t lino = 0;
6,703,004✔
3225

3226
  pCond->order = TSDB_ORDER_ASC;
6,703,004✔
3227
  pCond->numOfCols = (int32_t)taosArrayGetSize(colArray);
6,703,004✔
3228

3229
  pCond->colList = taosMemoryCalloc(pCond->numOfCols, sizeof(SColumnInfo));
6,703,004✔
3230
  QUERY_CHECK_NULL(pCond->colList, code, lino, _return, terrno);
6,703,004✔
3231

3232
  pCond->pSlotList = taosMemoryMalloc(sizeof(int32_t) * pCond->numOfCols);
6,703,004✔
3233
  QUERY_CHECK_NULL(pCond->pSlotList, code, lino, _return, terrno);
6,703,004✔
3234

3235
  pCond->twindows = pOrgCond->twindows;
6,703,004✔
3236
  pCond->type = pOrgCond->type;
6,703,004✔
3237
  pCond->startVersion = -1;
6,703,004✔
3238
  pCond->endVersion = -1;
6,703,004✔
3239
  pCond->skipRollup = true;
6,703,004✔
3240
  pCond->notLoadData = false;
6,703,004✔
3241

3242
  for (int32_t i = 0; i < pCond->numOfCols; ++i) {
33,746,906✔
3243
    SColIdPair* pColPair = taosArrayGet(colArray, i);
27,043,902✔
3244
    QUERY_CHECK_NULL(pColPair, code, lino, _return, terrno);
27,043,902✔
3245

3246
    bool find = false;
27,043,902✔
3247
    for (int32_t j = 0; j < pOrgCond->numOfCols; ++j) {
170,300,826✔
3248
      if (pOrgCond->colList[j].colId == pColPair->vtbColId) {
170,300,826✔
3249
        pCond->colList[i].type = pOrgCond->colList[j].type;
27,043,902✔
3250
        pCond->colList[i].bytes = pOrgCond->colList[j].bytes;
27,043,902✔
3251
        pCond->colList[i].colId = pColPair->orgColId;
27,043,902✔
3252
        pCond->colList[i].pk = pOrgCond->colList[j].pk;
27,043,902✔
3253
        pCond->pSlotList[i] = i;
27,043,902✔
3254
        find = true;
27,043,902✔
3255
        break;
27,043,902✔
3256
      }
3257
    }
3258
    QUERY_CHECK_CONDITION(find, code, lino, _return, TSDB_CODE_NOT_FOUND);
27,043,902✔
3259
  }
3260

3261
  return code;
6,703,004✔
3262
_return:
×
3263
  qError("%s failed at line %d since %s", __func__, lino, tstrerror(terrno));
×
3264
  taosMemoryFreeClear(pCond->colList);
×
3265
  taosMemoryFreeClear(pCond->pSlotList);
×
3266
  return code;
×
3267
}
3268

3269
void cleanupQueryTableDataCond(SQueryTableDataCond* pCond) {
450,289,652✔
3270
  taosMemoryFreeClear(pCond->colList);
450,289,652✔
3271
  taosMemoryFreeClear(pCond->pSlotList);
450,329,226✔
3272
}
450,307,970✔
3273

3274
int32_t convertFillType(int32_t mode) {
2,057,584✔
3275
  int32_t type = TSDB_FILL_NONE;
2,057,584✔
3276
  switch (mode) {
2,057,584✔
3277
    case FILL_MODE_PREV:
109,570✔
3278
      type = TSDB_FILL_PREV;
109,570✔
3279
      break;
109,570✔
3280
    case FILL_MODE_NONE:
×
3281
      type = TSDB_FILL_NONE;
×
3282
      break;
×
3283
    case FILL_MODE_NULL:
130,540✔
3284
      type = TSDB_FILL_NULL;
130,540✔
3285
      break;
130,540✔
3286
    case FILL_MODE_NULL_F:
11,334✔
3287
      type = TSDB_FILL_NULL_F;
11,334✔
3288
      break;
11,334✔
3289
    case FILL_MODE_NEXT:
96,782✔
3290
      type = TSDB_FILL_NEXT;
96,782✔
3291
      break;
96,782✔
3292
    case FILL_MODE_VALUE:
149,120✔
3293
      type = TSDB_FILL_SET_VALUE;
149,120✔
3294
      break;
149,120✔
3295
    case FILL_MODE_VALUE_F:
4,284✔
3296
      type = TSDB_FILL_SET_VALUE_F;
4,284✔
3297
      break;
4,284✔
3298
    case FILL_MODE_LINEAR:
152,853✔
3299
      type = TSDB_FILL_LINEAR;
152,853✔
3300
      break;
152,853✔
3301
    case FILL_MODE_NEAR:
1,403,101✔
3302
      type = TSDB_FILL_NEAR;
1,403,101✔
3303
      break;
1,403,101✔
3304
    default:
×
3305
      type = TSDB_FILL_NONE;
×
3306
  }
3307

3308
  return type;
2,057,584✔
3309
}
3310

3311
void getInitialStartTimeWindow(SInterval* pInterval, TSKEY ts, STimeWindow* w, bool ascQuery) {
1,802,391,138✔
3312
  if (ascQuery) {
1,802,391,138✔
3313
    *w = getAlignQueryTimeWindow(pInterval, ts);
1,802,843,603✔
3314
  } else {
3315
    // the start position of the first time window in the endpoint that spreads beyond the queried last timestamp
3316
    *w = getAlignQueryTimeWindow(pInterval, ts);
60,133✔
3317

3318
    int64_t key = w->skey;
146,428✔
3319
    while (key < ts) {  // moving towards end
160,953✔
3320
      key = getNextTimeWindowStart(pInterval, key, TSDB_ORDER_ASC);
76,080✔
3321
      if (key > ts) {
76,080✔
3322
        break;
61,555✔
3323
      }
3324

3325
      w->skey = key;
14,525✔
3326
    }
3327
    w->ekey = taosTimeAdd(w->skey, pInterval->interval, pInterval->intervalUnit, pInterval->precision, NULL) - 1;
146,428✔
3328
  }
3329
}
1,803,642,465✔
3330

3331
static STimeWindow doCalculateTimeWindow(int64_t ts, SInterval* pInterval) {
26,715,866✔
3332
  STimeWindow w = {0};
26,715,866✔
3333

3334
  w.skey = taosTimeTruncate(ts, pInterval);
26,715,866✔
3335
  w.ekey = taosTimeGetIntervalEnd(w.skey, pInterval);
26,714,873✔
3336
  return w;
26,718,648✔
3337
}
3338

3339
STimeWindow getFirstQualifiedTimeWindow(int64_t ts, STimeWindow* pWindow, SInterval* pInterval, int32_t order) {
1,631,792✔
3340
  STimeWindow win = *pWindow;
1,631,792✔
3341
  STimeWindow save = win;
1,631,792✔
3342
  while (win.skey <= ts && win.ekey >= ts) {
9,056,888✔
3343
    save = win;
7,425,096✔
3344
    // get previous time window
3345
    getNextTimeWindow(pInterval, &win, order == TSDB_ORDER_DESC ? TSDB_ORDER_ASC : TSDB_ORDER_DESC);
7,425,096✔
3346
  }
3347

3348
  return save;
1,631,792✔
3349
}
3350

3351
// get the correct time window according to the handled timestamp
3352
// todo refactor
3353
STimeWindow getActiveTimeWindow(SDiskbasedBuf* pBuf, SResultRowInfo* pResultRowInfo, int64_t ts, SInterval* pInterval,
45,448,017✔
3354
                                int32_t order) {
3355
  STimeWindow w = {0};
45,448,017✔
3356
  if (pResultRowInfo->cur.pageId == -1) {  // the first window, from the previous stored value
45,450,666✔
3357
    getInitialStartTimeWindow(pInterval, ts, &w, (order != TSDB_ORDER_DESC));
5,493,512✔
3358
    return w;
5,494,213✔
3359
  }
3360

3361
  SResultRow* pRow = getResultRowByPos(pBuf, &pResultRowInfo->cur, false);
39,959,951✔
3362
  if (pRow) {
39,961,577✔
3363
    TAOS_SET_OBJ_ALIGNED(&w, pRow->win);
39,961,946✔
3364
  }
3365

3366
  // in case of typical time window, we can calculate time window directly.
3367
  if (w.skey > ts || w.ekey < ts) {
39,962,403✔
3368
    w = doCalculateTimeWindow(ts, pInterval);
26,718,623✔
3369
  }
3370

3371
  if (pInterval->interval != pInterval->sliding) {
39,962,600✔
3372
    // it is an sliding window query, in which sliding value is not equalled to
3373
    // interval value, and we need to find the first qualified time window.
3374
    w = getFirstQualifiedTimeWindow(ts, &w, pInterval, order);
1,631,792✔
3375
  }
3376

3377
  return w;
39,960,772✔
3378
}
3379

3380
TSKEY getNextTimeWindowStart(const SInterval* pInterval, TSKEY start, int32_t order) {
2,147,483,647✔
3381
  int32_t factor = GET_FORWARD_DIRECTION_FACTOR(order);
2,147,483,647✔
3382
  TSKEY   nextStart = taosTimeAdd(start, -1 * pInterval->offset, pInterval->offsetUnit, pInterval->precision, NULL);
2,147,483,647✔
3383
  nextStart = taosTimeAdd(nextStart, factor * pInterval->sliding, pInterval->slidingUnit, pInterval->precision, NULL);
2,147,483,647✔
3384
  nextStart = taosTimeAdd(nextStart, pInterval->offset, pInterval->offsetUnit, pInterval->precision, NULL);
2,147,483,647✔
3385
  return nextStart;
2,147,483,647✔
3386
}
3387

3388
void getNextTimeWindow(const SInterval* pInterval, STimeWindow* tw, int32_t order) {
2,147,483,647✔
3389
  tw->skey = getNextTimeWindowStart(pInterval, tw->skey, order);
2,147,483,647✔
3390
  tw->ekey = taosTimeAdd(tw->skey, pInterval->interval, pInterval->intervalUnit, pInterval->precision, NULL) - 1;
2,147,483,647✔
3391
}
2,147,483,647✔
3392

3393
bool hasLimitOffsetInfo(SLimitInfo* pLimitInfo) {
310,103,031✔
3394
  return (pLimitInfo->limit.limit != -1 || pLimitInfo->limit.offset != -1 || pLimitInfo->slimit.limit != -1 ||
617,702,360✔
3395
          pLimitInfo->slimit.offset != -1);
307,598,087✔
3396
}
3397

3398
bool hasSlimitOffsetInfo(SLimitInfo* pLimitInfo) {
×
3399
  return (pLimitInfo->slimit.limit != -1 || pLimitInfo->slimit.offset != -1);
×
3400
}
3401

3402
void initLimitInfo(const SNode* pLimit, const SNode* pSLimit, SLimitInfo* pLimitInfo) {
463,498,464✔
3403
  SLimit limit = {.limit = getLimit(pLimit), .offset = getOffset(pLimit)};
463,498,464✔
3404
  SLimit slimit = {.limit = getLimit(pSLimit), .offset = getOffset(pSLimit)};
463,412,694✔
3405

3406
  pLimitInfo->limit = limit;
463,467,332✔
3407
  pLimitInfo->slimit = slimit;
463,473,116✔
3408
  pLimitInfo->remainOffset = limit.offset;
463,433,214✔
3409
  pLimitInfo->remainGroupOffset = slimit.offset;
463,397,765✔
3410
  pLimitInfo->numOfOutputRows = 0;
463,427,346✔
3411
  pLimitInfo->numOfOutputGroups = 0;
463,455,713✔
3412
  pLimitInfo->currentGroupId = 0;
463,436,966✔
3413
}
463,478,918✔
3414

3415
void resetLimitInfoForNextGroup(SLimitInfo* pLimitInfo) {
51,612,953✔
3416
  pLimitInfo->numOfOutputRows = 0;
51,612,953✔
3417
  pLimitInfo->remainOffset = pLimitInfo->limit.offset;
51,618,434✔
3418
}
51,618,652✔
3419

3420
int32_t tableListGetSize(const STableListInfo* pTableList, int32_t* pRes) {
470,366,995✔
3421
  if (taosArrayGetSize(pTableList->pTableList) != taosHashGetSize(pTableList->map)) {
470,366,995✔
3422
    qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(TSDB_CODE_QRY_EXECUTOR_INTERNAL_ERROR));
×
3423
    return TSDB_CODE_QRY_EXECUTOR_INTERNAL_ERROR;
×
3424
  }
3425
  (*pRes) = taosArrayGetSize(pTableList->pTableList);
470,440,450✔
3426
  return TSDB_CODE_SUCCESS;
470,437,857✔
3427
}
3428

3429
uint64_t tableListGetSuid(const STableListInfo* pTableList) { return pTableList->idInfo.suid; }
3,479,790✔
3430

3431
STableKeyInfo* tableListGetInfo(const STableListInfo* pTableList, int32_t index) {
154,102,251✔
3432
  if (taosArrayGetSize(pTableList->pTableList) == 0) {
154,102,251✔
3433
    return NULL;
3,648✔
3434
  }
3435

3436
  return taosArrayGet(pTableList->pTableList, index);
154,099,036✔
3437
}
3438

3439
int32_t tableListFind(const STableListInfo* pTableList, uint64_t uid, int32_t startIndex) {
9,884✔
3440
  int32_t numOfTables = taosArrayGetSize(pTableList->pTableList);
9,884✔
3441
  if (startIndex >= numOfTables) {
9,884✔
3442
    return -1;
×
3443
  }
3444

3445
  for (int32_t i = startIndex; i < numOfTables; ++i) {
107,718✔
3446
    STableKeyInfo* p = taosArrayGet(pTableList->pTableList, i);
107,718✔
3447
    if (!p) {
107,718✔
3448
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
3449
      return -1;
×
3450
    }
3451
    if (p->uid == uid) {
107,718✔
3452
      return i;
9,884✔
3453
    }
3454
  }
3455
  return -1;
×
3456
}
3457

3458
void tableListGetSourceTableInfo(const STableListInfo* pTableList, uint64_t* psuid, uint64_t* uid, int32_t* type) {
66,509✔
3459
  *psuid = pTableList->idInfo.suid;
66,509✔
3460
  *uid = pTableList->idInfo.uid;
66,509✔
3461
  *type = pTableList->idInfo.tableType;
66,509✔
3462
}
66,509✔
3463

3464
uint64_t tableListGetTableGroupId(const STableListInfo* pTableList, uint64_t tableUid) {
570,784,326✔
3465
  int32_t* slot = taosHashGet(pTableList->map, &tableUid, sizeof(tableUid));
570,784,326✔
3466
  if (slot == NULL) {
571,016,722✔
3467
    qDebug("table:%" PRIu64 " not found in table list", tableUid);
×
3468
    return -1;
×
3469
  }
3470

3471
  STableKeyInfo* pKeyInfo = taosArrayGet(pTableList->pTableList, *slot);
571,016,722✔
3472
  if (pKeyInfo == NULL) {
571,176,808✔
3473
    qDebug("table:%" PRIu64 " not found in table list", tableUid);
×
3474
    return -1;
×
3475
  }
3476
  return pKeyInfo->groupId;
571,176,808✔
3477
}
3478

3479
// TODO handle the group offset info, fix it, the rule of group output will be broken by this function
3480
// int32_t tableListRemoveTableInfo(STableListInfo* pTableList, uint64_t uid) {
3481
//   int32_t code = TSDB_CODE_SUCCESS;
3482
//   int32_t lino = 0;
3483

3484
//   int32_t* slot = taosHashGet(pTableList->map, &uid, sizeof(uid));
3485
//   if (slot == NULL) {
3486
//     qDebug("table:%" PRIu64 " not found in table list", uid);
3487
//     return 0;
3488
//   }
3489

3490
//   taosArrayRemove(pTableList->pTableList, *slot);
3491
//   code = taosHashRemove(pTableList->map, &uid, sizeof(uid));
3492

3493
//   _end:
3494
//   if (code != TSDB_CODE_SUCCESS) {
3495
//     qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
3496
//   } else {
3497
//     qDebug("uid:%" PRIu64 ", remove from table list", uid);
3498
//   }
3499

3500
//   return code;
3501
// }
3502

3503
int32_t tableListAddTableInfo(STableListInfo* pTableList, uint64_t uid, uint64_t gid) {
200,386✔
3504
  int32_t code = TSDB_CODE_SUCCESS;
200,386✔
3505
  int32_t lino = 0;
200,386✔
3506
  if (pTableList->map == NULL) {
200,386✔
3507
    pTableList->map = taosHashInit(32, taosGetDefaultHashFunction(TSDB_DATA_TYPE_BINARY), false, HASH_ENTRY_LOCK);
×
3508
    QUERY_CHECK_NULL(pTableList->map, code, lino, _end, terrno);
×
3509
  }
3510

3511
  STableKeyInfo keyInfo = {.uid = uid, .groupId = gid};
200,609✔
3512
  void*         p = taosHashGet(pTableList->map, &uid, sizeof(uid));
200,676✔
3513
  if (p != NULL) {
200,493✔
3514
    qInfo("table:%" PRId64 " already in tableIdList, ignore it", uid);
178✔
3515
    goto _end;
178✔
3516
  }
3517

3518
  void* tmp = taosArrayPush(pTableList->pTableList, &keyInfo);
200,315✔
3519
  QUERY_CHECK_NULL(tmp, code, lino, _end, terrno);
200,672✔
3520

3521
  int32_t slot = (int32_t)taosArrayGetSize(pTableList->pTableList) - 1;
200,672✔
3522
  code = taosHashPut(pTableList->map, &uid, sizeof(uid), &slot, sizeof(slot));
200,373✔
3523
  if (code != TSDB_CODE_SUCCESS) {
200,721✔
3524
    // we have checked the existence of uid in hash map above
3525
    QUERY_CHECK_CONDITION((code != TSDB_CODE_DUP_KEY), code, lino, _end, TSDB_CODE_QRY_EXECUTOR_INTERNAL_ERROR);
×
3526
    taosArrayPopTailBatch(pTableList->pTableList, 1);  // let's pop the last element in the array list
×
3527
  }
3528

3529
_end:
200,899✔
3530
  if (code != TSDB_CODE_SUCCESS) {
200,667✔
3531
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
3532
  } else {
3533
    qDebug("uid:%" PRIu64 ", groupId:%" PRIu64 " added into table list, slot:%d, total:%d", uid, gid, slot, slot + 1);
200,667✔
3534
  }
3535

3536
  return code;
200,899✔
3537
}
3538

3539
int32_t tableListGetGroupList(const STableListInfo* pTableList, int32_t ordinalGroupIndex, STableKeyInfo** pKeyInfo,
186,107,014✔
3540
                              int32_t* size) {
3541
  int32_t totalGroups = tableListGetOutputGroups(pTableList);
186,107,014✔
3542
  int32_t numOfTables = 0;
186,170,448✔
3543
  int32_t code = tableListGetSize(pTableList, &numOfTables);
186,177,155✔
3544
  if (code != TSDB_CODE_SUCCESS) {
186,217,227✔
3545
    qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(code));
×
3546
    return code;
×
3547
  }
3548

3549
  if (ordinalGroupIndex < 0 || ordinalGroupIndex >= totalGroups) {
186,217,227✔
3550
    return TSDB_CODE_INVALID_PARA;
×
3551
  }
3552

3553
  // here handle two special cases:
3554
  // 1. only one group exists, and 2. one table exists for each group.
3555
  if (totalGroups == 1) {
186,217,227✔
3556
    *size = numOfTables;
185,823,787✔
3557
    *pKeyInfo = (*size == 0) ? NULL : taosArrayGet(pTableList->pTableList, 0);
185,894,369✔
3558
    return TSDB_CODE_SUCCESS;
185,858,751✔
3559
  } else if (totalGroups == numOfTables) {
393,440✔
3560
    *size = 1;
331,634✔
3561
    *pKeyInfo = taosArrayGet(pTableList->pTableList, ordinalGroupIndex);
331,634✔
3562
    return TSDB_CODE_SUCCESS;
331,634✔
3563
  }
3564

3565
  int32_t offset = pTableList->groupOffset[ordinalGroupIndex];
61,806✔
3566
  if (ordinalGroupIndex < totalGroups - 1) {
54,616✔
3567
    *size = pTableList->groupOffset[ordinalGroupIndex + 1] - offset;
40,111✔
3568
  } else {
3569
    *size = numOfTables - offset;
14,505✔
3570
  }
3571

3572
  *pKeyInfo = taosArrayGet(pTableList->pTableList, offset);
54,616✔
3573
  return TSDB_CODE_SUCCESS;
54,102✔
3574
}
3575

3576
int32_t tableListGetOutputGroups(const STableListInfo* pTableList) { return pTableList->numOfOuputGroups; }
551,230,786✔
3577

3578
bool oneTableForEachGroup(const STableListInfo* pTableList) { return pTableList->oneTableForEachGroup; }
661,031✔
3579

3580
STableListInfo* tableListCreate() {
219,803,276✔
3581
  STableListInfo* pListInfo = taosMemoryCalloc(1, sizeof(STableListInfo));
219,803,276✔
3582
  if (pListInfo == NULL) {
219,739,234✔
3583
    return NULL;
×
3584
  }
3585

3586
  pListInfo->remainGroups = NULL;
219,739,234✔
3587
  pListInfo->pTableList = taosArrayInit(4, sizeof(STableKeyInfo));
219,745,766✔
3588
  if (pListInfo->pTableList == NULL) {
219,820,674✔
3589
    goto _error;
×
3590
  }
3591

3592
  pListInfo->map = taosHashInit(1024, taosGetDefaultHashFunction(TSDB_DATA_TYPE_BIGINT), false, HASH_ENTRY_LOCK);
219,804,605✔
3593
  if (pListInfo->map == NULL) {
219,889,033✔
3594
    goto _error;
×
3595
  }
3596

3597
  pListInfo->numOfOuputGroups = 1;
219,879,230✔
3598
  return pListInfo;
219,886,139✔
3599

3600
_error:
×
3601
  tableListDestroy(pListInfo);
×
3602
  return NULL;
×
3603
}
3604

3605
void tableListDestroy(STableListInfo* pTableListInfo) {
228,970,024✔
3606
  if (pTableListInfo == NULL) {
228,970,024✔
3607
    return;
9,181,801✔
3608
  }
3609

3610
  taosArrayDestroy(pTableListInfo->pTableList);
219,788,223✔
3611
  taosMemoryFreeClear(pTableListInfo->groupOffset);
219,808,407✔
3612

3613
  taosHashCleanup(pTableListInfo->map);
219,810,207✔
3614
  taosHashCleanup(pTableListInfo->remainGroups);
219,845,750✔
3615
  pTableListInfo->pTableList = NULL;
219,876,618✔
3616
  pTableListInfo->map = NULL;
219,870,347✔
3617
  taosMemoryFree(pTableListInfo);
219,843,908✔
3618
}
3619

3620
void tableListClear(STableListInfo* pTableListInfo) {
149,468✔
3621
  if (pTableListInfo == NULL) {
149,468✔
3622
    return;
×
3623
  }
3624

3625
  taosArrayClear(pTableListInfo->pTableList);
149,468✔
3626
  taosHashClear(pTableListInfo->map);
149,725✔
3627
  taosHashClear(pTableListInfo->remainGroups);
149,874✔
3628
  taosMemoryFree(pTableListInfo->groupOffset);
149,874✔
3629
  pTableListInfo->numOfOuputGroups = 1;
149,816✔
3630
  pTableListInfo->oneTableForEachGroup = false;
149,874✔
3631
}
3632

3633
static int32_t orderbyGroupIdComparFn(const void* p1, const void* p2) {
498,937,769✔
3634
  STableKeyInfo* pInfo1 = (STableKeyInfo*)p1;
498,937,769✔
3635
  STableKeyInfo* pInfo2 = (STableKeyInfo*)p2;
498,937,769✔
3636

3637
  if (pInfo1->groupId == pInfo2->groupId) {
498,937,769✔
3638
    return 0;
482,722,583✔
3639
  } else {
3640
    return pInfo1->groupId < pInfo2->groupId ? -1 : 1;
16,218,059✔
3641
  }
3642
}
3643

3644
int32_t sortTableGroup(STableListInfo* pTableListInfo) {
18,161,055✔
3645
  int32_t code = TSDB_CODE_SUCCESS;
18,161,055✔
3646
  taosArraySort(pTableListInfo->pTableList, orderbyGroupIdComparFn);
18,161,055✔
3647
  int32_t size = taosArrayGetSize(pTableListInfo->pTableList);
18,166,662✔
3648
  if (size == 0) {
18,164,462✔
3649
    pTableListInfo->numOfOuputGroups = 0;
×
3650
    return code;
×
3651
  }
3652

3653
  SArray* pList = taosArrayInit(4, sizeof(int32_t));
18,164,462✔
3654
  if (!pList) {
18,164,814✔
3655
    code = terrno;
×
3656
    goto end;
×
3657
  }
3658

3659
  STableKeyInfo* pInfo = taosArrayGet(pTableListInfo->pTableList, 0);
18,164,814✔
3660
  if (pInfo == NULL) {
18,158,007✔
3661
    qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
3662
    code = terrno;
×
3663
    goto end;
×
3664
  }
3665
  uint64_t gid = pInfo->groupId;
18,158,007✔
3666

3667
  int32_t start = 0;
18,159,142✔
3668
  void*   tmp = taosArrayPush(pList, &start);
18,166,850✔
3669
  if (tmp == NULL) {
18,166,850✔
3670
    qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
3671
    code = terrno;
×
3672
    goto end;
×
3673
  }
3674

3675
  for (int32_t i = 1; i < size; ++i) {
119,912,113✔
3676
    pInfo = taosArrayGet(pTableListInfo->pTableList, i);
101,746,421✔
3677
    if (pInfo == NULL) {
101,745,053✔
3678
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
3679
      code = terrno;
×
3680
      goto end;
×
3681
    }
3682
    if (pInfo->groupId != gid) {
101,745,053✔
3683
      tmp = taosArrayPush(pList, &i);
3,535,822✔
3684
      if (tmp == NULL) {
3,535,822✔
3685
        qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
3686
        code = terrno;
×
3687
        goto end;
×
3688
      }
3689
      gid = pInfo->groupId;
3,535,822✔
3690
    }
3691
  }
3692

3693
  pTableListInfo->numOfOuputGroups = taosArrayGetSize(pList);
18,164,047✔
3694
  pTableListInfo->groupOffset = taosMemoryMalloc(sizeof(int32_t) * pTableListInfo->numOfOuputGroups);
18,164,467✔
3695
  if (pTableListInfo->groupOffset == NULL) {
18,161,861✔
3696
    code = terrno;
×
3697
    goto end;
×
3698
  }
3699

3700
  memcpy(pTableListInfo->groupOffset, taosArrayGet(pList, 0), sizeof(int32_t) * pTableListInfo->numOfOuputGroups);
18,156,696✔
3701

3702
end:
18,163,795✔
3703
  taosArrayDestroy(pList);
18,160,618✔
3704
  return code;
18,160,523✔
3705
}
3706

3707
int32_t buildGroupIdMapForAllTables(STableListInfo* pTableListInfo, SReadHandle* pHandle, SScanPhysiNode* pScanNode,
202,237,031✔
3708
                                    SNodeList* group, bool groupSort, uint8_t* digest, SStorageAPI* pAPI, SHashObj* groupIdMap) {
3709
  int32_t code = TSDB_CODE_SUCCESS;
202,237,031✔
3710

3711
  bool   groupByTbname = groupbyTbname(group);
202,237,031✔
3712
  size_t numOfTables = taosArrayGetSize(pTableListInfo->pTableList);
202,243,969✔
3713
  if (!numOfTables) {
202,230,273✔
3714
    return code;
6,214✔
3715
  }
3716
  qDebug("numOfTables:%zu, groupByTbname:%d, group:%p", numOfTables, groupByTbname, group);
202,224,059✔
3717
  if (group == NULL || groupByTbname) {
202,198,229✔
3718
    if (tsCountAlwaysReturnValue && QUERY_NODE_PHYSICAL_PLAN_TABLE_SCAN == nodeType(pScanNode) &&
200,294,791✔
3719
        ((STableScanPhysiNode*)pScanNode)->needCountEmptyTable) {
173,438,160✔
3720
      pTableListInfo->remainGroups =
7,253,198✔
3721
          taosHashInit(numOfTables, taosGetDefaultHashFunction(TSDB_DATA_TYPE_BIGINT), false, HASH_NO_LOCK);
7,253,198✔
3722
      if (pTableListInfo->remainGroups == NULL) {
7,253,198✔
3723
        return terrno;
×
3724
      }
3725

3726
      for (int i = 0; i < numOfTables; i++) {
26,060,827✔
3727
        STableKeyInfo* info = taosArrayGet(pTableListInfo->pTableList, i);
18,807,784✔
3728
        if (!info) {
18,807,068✔
3729
          qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
3730
          return terrno;
×
3731
        }
3732
        info->groupId = groupByTbname ? info->uid : 0;
18,807,068✔
3733
        int32_t tempRes = taosHashPut(pTableListInfo->remainGroups, &(info->groupId), sizeof(info->groupId),
18,807,939✔
3734
                                      &(info->uid), sizeof(info->uid));
18,808,094✔
3735
        if (tempRes != TSDB_CODE_SUCCESS && tempRes != TSDB_CODE_DUP_KEY) {
18,807,629✔
3736
          qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(tempRes));
×
3737
          return tempRes;
×
3738
        }
3739
      }
3740
    } else {
3741
      for (int32_t i = 0; i < numOfTables; i++) {
658,121,024✔
3742
        STableKeyInfo* info = taosArrayGet(pTableListInfo->pTableList, i);
465,065,980✔
3743
        if (!info) {
465,010,711✔
3744
          qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
3745
          return terrno;
×
3746
        }
3747
        info->groupId = groupByTbname ? info->uid : 0;
465,010,711✔
3748
        
3749
      }
3750
    }
3751
    if (groupIdMap && group != NULL){
200,308,087✔
3752
      getColInfoResultForGroupbyForStream(pHandle->vnode, group, pTableListInfo, pAPI, groupIdMap);
107,663✔
3753
    }
3754

3755
    pTableListInfo->oneTableForEachGroup = groupByTbname;
200,308,087✔
3756
    if (numOfTables == 1 && pTableListInfo->idInfo.tableType == TSDB_CHILD_TABLE) {
200,271,198✔
3757
      pTableListInfo->oneTableForEachGroup = true;
85,679,976✔
3758
    }
3759

3760
    if (groupSort && groupByTbname) {
200,222,681✔
3761
      taosArraySort(pTableListInfo->pTableList, orderbyGroupIdComparFn);
1,461,520✔
3762
      pTableListInfo->numOfOuputGroups = numOfTables;
1,461,602✔
3763
    } else if (groupByTbname && pScanNode->groupOrderScan) {
198,761,161✔
3764
      pTableListInfo->numOfOuputGroups = numOfTables;
32,225✔
3765
    } else {
3766
      pTableListInfo->numOfOuputGroups = 1;
198,729,046✔
3767
    }
3768
    if (groupSort || pScanNode->groupOrderScan) {
200,316,042✔
3769
      code = sortTableGroup(pTableListInfo);
18,012,288✔
3770
    }
3771
  } else {
3772
    bool initRemainGroups = false;
1,903,438✔
3773
    if (QUERY_NODE_PHYSICAL_PLAN_TABLE_SCAN == nodeType(pScanNode)) {
1,903,438✔
3774
      STableScanPhysiNode* pTableScanNode = (STableScanPhysiNode*)pScanNode;
1,706,825✔
3775
      if (tsCountAlwaysReturnValue && pTableScanNode->needCountEmptyTable &&
1,706,825✔
3776
          !(groupSort || pScanNode->groupOrderScan)) {
882,929✔
3777
        initRemainGroups = true;
856,031✔
3778
      }
3779
    }
3780

3781
    code = getColInfoResultForGroupby(pHandle->vnode, group, pTableListInfo, digest, pAPI, initRemainGroups, groupIdMap);
1,903,438✔
3782
    if (code != TSDB_CODE_SUCCESS) {
1,902,750✔
3783
      return code;
×
3784
    }
3785

3786
    if (pScanNode->groupOrderScan) pTableListInfo->numOfOuputGroups = taosArrayGetSize(pTableListInfo->pTableList);
1,902,750✔
3787

3788
    if (groupSort || pScanNode->groupOrderScan) {
1,903,438✔
3789
      code = sortTableGroup(pTableListInfo);
152,888✔
3790
    }
3791
  }
3792

3793
  // add all table entry in the hash map
3794
  size_t size = taosArrayGetSize(pTableListInfo->pTableList);
202,214,991✔
3795
  for (int32_t i = 0; i < size; ++i) {
697,202,335✔
3796
    STableKeyInfo* p = taosArrayGet(pTableListInfo->pTableList, i);
494,924,190✔
3797
    if (!p) {
494,856,808✔
3798
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
3799
      return terrno;
×
3800
    }
3801
    int32_t tempRes = taosHashPut(pTableListInfo->map, &p->uid, sizeof(uint64_t), &i, sizeof(int32_t));
494,856,808✔
3802
    if (tempRes != TSDB_CODE_SUCCESS && tempRes != TSDB_CODE_DUP_KEY) {
495,028,339✔
3803
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(tempRes));
×
3804
      return tempRes;
×
3805
    }
3806
  }
3807

3808
  return code;
202,290,162✔
3809
}
3810

3811
int32_t createScanTableListInfo(SScanPhysiNode* pScanNode, SNodeList* pGroupTags, bool groupSort, SReadHandle* pHandle,
214,147,425✔
3812
                                STableListInfo* pTableListInfo, SNode* pTagCond, SNode* pTagIndexCond,
3813
                                SExecTaskInfo* pTaskInfo, SHashObj* groupIdMap) {
3814
  int64_t     st = taosGetTimestampUs();
214,160,581✔
3815
  const char* idStr = GET_TASKID(pTaskInfo);
214,160,581✔
3816

3817
  if (pHandle == NULL) {
214,018,717✔
3818
    qError("invalid handle, in creating operator tree, %s", idStr);
×
3819
    return TSDB_CODE_INVALID_PARA;
×
3820
  }
3821

3822
  if (pHandle->uid != 0) {
214,018,717✔
3823
    pScanNode->uid = pHandle->uid;
46,404✔
3824
    pScanNode->tableType = TSDB_CHILD_TABLE;
46,404✔
3825
  }
3826
  uint8_t digest[17] = {0};
214,030,315✔
3827
  int32_t code = getTableList(pHandle->vnode, pScanNode, pTagCond, pTagIndexCond, pTableListInfo, digest, idStr,
214,141,278✔
3828
                              &pTaskInfo->storageAPI, pTaskInfo->pStreamRuntimeInfo);
214,113,261✔
3829
  if (code != TSDB_CODE_SUCCESS) {
214,205,174✔
3830
    qError("failed to getTableList, code:%s", tstrerror(code));
1,100✔
3831
    return code;
1,100✔
3832
  }
3833

3834
  int32_t numOfTables = taosArrayGetSize(pTableListInfo->pTableList);
214,204,074✔
3835

3836
  int64_t st1 = taosGetTimestampUs();
214,226,843✔
3837
  pTaskInfo->cost.extractListTime = (st1 - st) / 1000.0;
214,226,843✔
3838
  qDebug("extract queried table list completed, %d tables, elapsed time:%.2f ms %s", numOfTables,
214,275,404✔
3839
         pTaskInfo->cost.extractListTime, idStr);
3840

3841
  if (numOfTables == 0) {
214,270,671✔
3842
    qDebug("no table qualified for query, %s", idStr);
12,084,917✔
3843
    return TSDB_CODE_SUCCESS;
12,084,917✔
3844
  }
3845

3846
  code = buildGroupIdMapForAllTables(pTableListInfo, pHandle, pScanNode, pGroupTags, groupSort, digest, &pTaskInfo->storageAPI, groupIdMap);
202,185,754✔
3847
  if (code != TSDB_CODE_SUCCESS) {
202,167,367✔
3848
    return code;
×
3849
  }
3850

3851
  pTaskInfo->cost.groupIdMapTime = (taosGetTimestampUs() - st1) / 1000.0;
202,182,229✔
3852
  qDebug("generate group id map completed, elapsed time:%.2f ms %s", pTaskInfo->cost.groupIdMapTime, idStr);
202,149,328✔
3853

3854
  return TSDB_CODE_SUCCESS;
202,155,892✔
3855
}
3856

3857
char* getStreamOpName(uint16_t opType) {
9,793,535✔
3858
  switch (opType) {
9,793,535✔
3859
    case QUERY_NODE_PHYSICAL_PLAN_STREAM_SCAN:
×
3860
      return "stream scan";
×
3861
    case QUERY_NODE_PHYSICAL_PLAN_PROJECT:
9,565,077✔
3862
      return "project";
9,565,077✔
3863
    case QUERY_NODE_PHYSICAL_PLAN_EXTERNAL_WINDOW:
228,458✔
3864
      return "external window";
228,458✔
3865
  }
3866
  return "error name";
×
3867
}
3868

3869
void printDataBlock(SSDataBlock* pBlock, const char* flag, const char* taskIdStr, int64_t qId) {
266,706,553✔
3870
  if (qDebugFlag & DEBUG_TRACE) {
266,706,553✔
3871
    if (!pBlock) {
38,352✔
3872
      qDebug("%" PRIx64 " %s %s %s: Block is Null", qId, taskIdStr, flag, __func__);
6,936✔
3873
      return;
6,936✔
3874
    } else if (pBlock->info.rows == 0) {
31,416✔
3875
      qDebug("%" PRIx64 " %s %s %s: Block is Empty. block type %d", qId, taskIdStr, flag, __func__, pBlock->info.type);
×
3876
      return;
×
3877
    }
3878
    
3879
    char*   pBuf = NULL;
31,416✔
3880
    int32_t code = dumpBlockData(pBlock, flag, &pBuf, taskIdStr, qId);
31,416✔
3881
    if (code == 0) {
31,416✔
3882
      qDebugL("%" PRIx64 " %s %s", qId, __func__, pBuf);
31,416✔
3883
      taosMemoryFree(pBuf);
31,416✔
3884
    }
3885
  }
3886
}
3887

3888
void printSpecDataBlock(SSDataBlock* pBlock, const char* flag, const char* opStr, const char* taskIdStr) {
×
3889
  if (!pBlock) {
×
3890
    qDebug("%s===stream===%s %s: Block is Null", taskIdStr, flag, opStr);
×
3891
    return;
×
3892
  } else if (pBlock->info.rows == 0) {
×
3893
    qDebug("%s===stream===%s %s: Block is Empty. block type %d.skey:%" PRId64 ",ekey:%" PRId64 ",version%" PRId64,
×
3894
           taskIdStr, flag, opStr, pBlock->info.type, pBlock->info.window.skey, pBlock->info.window.ekey,
3895
           pBlock->info.version);
3896
    return;
×
3897
  }
3898
  if (qDebugFlag & DEBUG_TRACE) {
×
3899
    char* pBuf = NULL;
×
3900
    char  flagBuf[64];
×
3901
    snprintf(flagBuf, sizeof(flagBuf), "%s %s", flag, opStr);
×
3902
    int32_t code = dumpBlockData(pBlock, flagBuf, &pBuf, taskIdStr, 0);
×
3903
    if (code == 0) {
×
3904
      qDebug("%s", pBuf);
×
3905
      taosMemoryFree(pBuf);
×
3906
    }
3907
  }
3908
}
3909

3910
TSKEY getStartTsKey(STimeWindow* win, const TSKEY* tsCols) { return tsCols == NULL ? win->skey : tsCols[0]; }
12,143,321✔
3911

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

3915
  int64_t duration = pWin->ekey > pWin->skey ? pWin->ekey - pWin->skey + delta : pWin->skey - pWin->ekey + delta;
2,147,483,647✔
3916
  ts[2] = duration;            // set the duration
2,147,483,647✔
3917
  ts[3] = pWin->skey;          // window start key
2,147,483,647✔
3918
  ts[4] = pWin->ekey + delta;  // window end key
2,147,483,647✔
3919
}
2,147,483,647✔
3920

3921
int32_t compKeys(const SArray* pSortGroupCols, const char* oldkeyBuf, int32_t oldKeysLen, const SSDataBlock* pBlock,
940,853,032✔
3922
                 int32_t rowIndex) {
3923
  SColumnDataAgg* pColAgg = NULL;
940,853,032✔
3924
  const char*     isNull = oldkeyBuf;
940,853,032✔
3925
  const char*     p = oldkeyBuf + sizeof(int8_t) * pSortGroupCols->size;
940,853,032✔
3926

3927
  for (int32_t i = 0; i < pSortGroupCols->size; ++i) {
2,147,483,647✔
3928
    const SColumn*         pCol = (SColumn*)TARRAY_GET_ELEM(pSortGroupCols, i);
1,457,503,363✔
3929
    const SColumnInfoData* pColInfoData = TARRAY_GET_ELEM(pBlock->pDataBlock, pCol->slotId);
1,459,709,916✔
3930
    if (pBlock->pBlockAgg) pColAgg = &pBlock->pBlockAgg[pCol->slotId];
1,459,702,316✔
3931

3932
    if (colDataIsNull(pColInfoData, pBlock->info.rows, rowIndex, pColAgg)) {
2,147,483,647✔
3933
      if (isNull[i] != 1) return 1;
100,051,150✔
3934
    } else {
3935
      if (isNull[i] != 0) return 1;
1,360,160,741✔
3936
      const char* val = colDataGetData(pColInfoData, rowIndex);
1,359,335,394✔
3937
      if (pCol->type == TSDB_DATA_TYPE_JSON) {
1,359,461,840✔
3938
        int32_t len = getJsonValueLen(val);
×
3939
        if (memcmp(p, val, len) != 0) return 1;
×
3940
        p += len;
×
3941
      } else if (IS_VAR_DATA_TYPE(pCol->type)) {
1,359,021,791✔
3942
        if (IS_STR_DATA_BLOB(pCol->type)) {
461,224,773✔
3943
          if (memcmp(p, val, blobDataTLen(val)) != 0) return 1;
×
3944
          p += blobDataTLen(val);
×
3945
        } else {
3946
          if (memcmp(p, val, varDataTLen(val)) != 0) return 1;
461,693,779✔
3947
          p += varDataTLen(val);
454,665,854✔
3948
        }
3949
      } else {
3950
        if (0 != memcmp(p, val, pCol->bytes)) return 1;
897,904,160✔
3951
        p += pCol->bytes;
880,691,176✔
3952
      }
3953
    }
3954
  }
3955
  if ((int32_t)(p - oldkeyBuf) != oldKeysLen) return 1;
916,180,623✔
3956
  return 0;
916,167,755✔
3957
}
3958

3959
int32_t buildKeys(char* keyBuf, const SArray* pSortGroupCols, const SSDataBlock* pBlock, int32_t rowIndex) {
24,967,029✔
3960
  uint32_t        colNum = pSortGroupCols->size;
24,967,029✔
3961
  SColumnDataAgg* pColAgg = NULL;
24,967,236✔
3962
  char*           isNull = keyBuf;
24,967,236✔
3963
  char*           p = keyBuf + sizeof(int8_t) * colNum;
24,967,236✔
3964

3965
  for (int32_t i = 0; i < colNum; ++i) {
74,554,025✔
3966
    const SColumn*         pCol = (SColumn*)TARRAY_GET_ELEM(pSortGroupCols, i);
49,586,168✔
3967
    const SColumnInfoData* pColInfoData = TARRAY_GET_ELEM(pBlock->pDataBlock, pCol->slotId);
49,625,084✔
3968
    if (pCol->slotId > pBlock->pDataBlock->size) continue;
49,625,084✔
3969

3970
    if (pBlock->pBlockAgg) pColAgg = &pBlock->pBlockAgg[pCol->slotId];
49,626,119✔
3971

3972
    if (colDataIsNull(pColInfoData, pBlock->info.rows, rowIndex, pColAgg)) {
99,206,698✔
3973
      isNull[i] = 1;
1,784,390✔
3974
    } else {
3975
      isNull[i] = 0;
47,844,834✔
3976
      const char* val = colDataGetData(pColInfoData, rowIndex);
47,796,396✔
3977
      if (pCol->type == TSDB_DATA_TYPE_JSON) {
47,839,452✔
3978
        int32_t len = getJsonValueLen(val);
×
3979
        memcpy(p, val, len);
×
3980
        p += len;
×
3981
      } else if (IS_VAR_DATA_TYPE(pCol->type)) {
47,840,487✔
3982
        if (IS_STR_DATA_BLOB(pCol->type)) {
7,178,333✔
3983
          blobDataCopy(p, val);
×
3984
          p += blobDataTLen(val);
×
3985
        } else {
3986
          varDataCopy(p, val);
7,196,963✔
3987
          p += varDataTLen(val);
7,196,756✔
3988
        }
3989
      } else {
3990
        memcpy(p, val, pCol->bytes);
40,644,766✔
3991
        p += pCol->bytes;
40,640,833✔
3992
      }
3993
    }
3994
  }
3995
  return (int32_t)(p - keyBuf);
24,967,857✔
3996
}
3997

3998
uint64_t calcGroupId(char* pData, int32_t len) {
2,147,483,647✔
3999
  T_MD5_CTX context;
2,147,483,647✔
4000
  tMD5Init(&context);
2,147,483,647✔
4001
  tMD5Update(&context, (uint8_t*)pData, len);
2,147,483,647✔
4002
  tMD5Final(&context);
2,147,483,647✔
4003

4004
  // NOTE: only extract the initial 8 bytes of the final MD5 digest
4005
  uint64_t id = 0;
2,147,483,647✔
4006
  memcpy(&id, context.digest, sizeof(uint64_t));
2,147,483,647✔
4007
  if (0 == id) memcpy(&id, context.digest + 8, sizeof(uint64_t));
2,147,483,647✔
4008
  return id;
2,147,483,647✔
4009
}
4010

4011
SNodeList* makeColsNodeArrFromSortKeys(SNodeList* pSortKeys) {
41,572✔
4012
  SNode*     node;
4013
  SNodeList* ret = NULL;
41,572✔
4014
  FOREACH(node, pSortKeys) {
126,640✔
4015
    SOrderByExprNode* pSortKey = (SOrderByExprNode*)node;
85,068✔
4016
    int32_t           code = nodesListMakeAppend(&ret, pSortKey->pExpr);
85,068✔
4017
    if (code != TSDB_CODE_SUCCESS) {
85,068✔
4018
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(code));
×
4019
      terrno = code;
×
4020
      return NULL;
×
4021
    }
4022
  }
4023
  return ret;
41,572✔
4024
}
4025

4026
int32_t extractKeysLen(const SArray* keys, int32_t* pLen) {
41,572✔
4027
  int32_t code = TSDB_CODE_SUCCESS;
41,572✔
4028
  int32_t lino = 0;
41,572✔
4029
  int32_t len = 0;
41,572✔
4030
  int32_t keyNum = taosArrayGetSize(keys);
41,572✔
4031
  for (int32_t i = 0; i < keyNum; ++i) {
106,000✔
4032
    SColumn* pCol = (SColumn*)taosArrayGet(keys, i);
64,428✔
4033
    QUERY_CHECK_NULL(pCol, code, lino, _end, terrno);
64,428✔
4034
    len += pCol->bytes;
64,428✔
4035
  }
4036
  len += sizeof(int8_t) * keyNum;  // null flag
41,572✔
4037
  *pLen = len;
41,572✔
4038

4039
_end:
41,572✔
4040
  if (code != TSDB_CODE_SUCCESS) {
41,572✔
4041
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
4042
  }
4043
  return code;
41,572✔
4044
}
4045

4046
int32_t parseErrorMsgFromAnalyticServer(SJson* pJson, const char* pId) {
×
4047
  int32_t code = TSDB_CODE_ANA_ANODE_RETURN_ERROR;
×
4048
  if (pJson == NULL) {
×
4049
    return code;
×
4050
  }
4051

4052
  char    pMsg[1024] = {0};
×
4053
  int32_t ret = tjsonGetStringValue(pJson, "msg", pMsg);
×
4054

4055
  if (ret == 0) {
×
4056
    qError("%s failed to exec imputation, msg:%s", pId, pMsg);
×
4057
    if (strstr(pMsg, "white noise") != NULL) {
×
4058
      code = TSDB_CODE_ANA_WN_DATA;
×
4059
    } else if (strstr(pMsg, "white-noise") != NULL) {
×
4060
      code = TSDB_CODE_ANA_WN_DATA;
×
4061
    } else if (strstr(pMsg, "[Errno 111] Connection refused") != NULL) {
×
4062
      code = TSDB_CODE_ANA_ALGO_NOT_LOAD;
×
4063
    }
4064
  } else {
4065
    qError("%s failed to extract msg from server, unknown error", pId);
×
4066
  }
4067

4068
  return code;
×
4069
}
4070

4071

4072
int32_t createBlockFromRemoteValueNode(SSDataBlock** ppBlock, SRemoteValueNode* pRemote) {
22,513,246✔
4073
  SValueNode* pVal = (SValueNode*)pRemote;
22,513,246✔
4074
  int32_t code = 0;
22,513,246✔
4075
  SSDataBlock* pBlock = taosMemoryCalloc(1, sizeof(SSDataBlock));
22,513,246✔
4076
  if (pBlock == NULL) {
22,511,808✔
4077
    return terrno;
×
4078
  }
4079

4080
  pBlock->pDataBlock = taosArrayInit(1, sizeof(SColumnInfoData));
22,511,808✔
4081
  if (pBlock->pDataBlock == NULL) {
22,511,808✔
4082
    code = terrno;
×
4083
    taosMemoryFree(pBlock);
×
4084
    return code;
×
4085
  }
4086

4087
  SColumnInfoData idata =
22,511,825✔
4088
      createColumnInfoData(pVal->node.resType.type, pVal->node.resType.bytes, 0);
22,512,362✔
4089
  idata.info.scale = pVal->node.resType.scale;
22,512,804✔
4090
  idata.info.precision = pVal->node.resType.precision;
22,513,246✔
4091

4092
  code = blockDataAppendColInfo(pBlock, &idata);
22,512,804✔
4093
  if (code != TSDB_CODE_SUCCESS) {
22,512,267✔
4094
    qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(code));
×
4095
    blockDataDestroy(pBlock);
×
4096
    *ppBlock = NULL;
×
4097
    return code;
×
4098
  }
4099

4100
  *ppBlock = pBlock;
22,512,267✔
4101

4102
  return code;
22,512,804✔
4103
}
4104

4105

4106
int32_t extractSingleRspBlock(SRetrieveTableRsp* pRetrieveRsp, SSDataBlock* pb) {
22,512,267✔
4107
  int32_t            code = TSDB_CODE_SUCCESS;
22,512,267✔
4108
  int32_t            lino = 0;
22,512,267✔
4109
  void*              decompBuf = NULL;
22,512,267✔
4110

4111
  char* pNextStart = pRetrieveRsp->data;
22,512,267✔
4112
  char* pStart = pNextStart;
22,512,205✔
4113

4114
  int32_t index = 0;
22,511,920✔
4115

4116
  if (pRetrieveRsp->compressed) {  // decompress the data
22,511,920✔
4117
    decompBuf = taosMemoryMalloc(pRetrieveRsp->payloadLen);
×
4118
    QUERY_CHECK_NULL(decompBuf, code, lino, _end, terrno);
×
4119
  }
4120

4121
  int32_t compLen = *(int32_t*)pStart;
22,512,205✔
4122
  pStart += sizeof(int32_t);
22,512,205✔
4123

4124
  int32_t rawLen = *(int32_t*)pStart;
22,513,246✔
4125
  pStart += sizeof(int32_t);
22,513,246✔
4126
  QUERY_CHECK_CONDITION((compLen <= rawLen && compLen != 0), code, lino, _end, TSDB_CODE_QRY_EXECUTOR_INTERNAL_ERROR);
22,512,804✔
4127

4128
  pNextStart = pStart + compLen;
22,512,804✔
4129
  if (pRetrieveRsp->compressed && (compLen < rawLen)) {
22,510,879✔
4130
    int32_t t = tsDecompressString(pStart, compLen, 1, decompBuf, rawLen, ONE_STAGE_COMP, NULL, 0);
×
4131
    QUERY_CHECK_CONDITION((t == rawLen), code, lino, _end, TSDB_CODE_QRY_EXECUTOR_INTERNAL_ERROR);
×
4132
    pStart = decompBuf;
×
4133
  }
4134

4135
  code = blockDecodeInternal(pb, pStart, (const char**)&pStart);
22,511,226✔
4136
  if (code != 0) {
22,509,805✔
4137
    taosMemoryFreeClear(pRetrieveRsp);
×
4138
    goto _end;
×
4139
  }
4140

4141
_end:
22,509,805✔
4142
  if (code != TSDB_CODE_SUCCESS) {
22,509,805✔
4143
    blockDataDestroy(pb);
×
4144
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
4145
  }
4146
  return code;
22,509,805✔
4147
}
4148

4149
int32_t setValueFromResBlock(STaskSubJobCtx* ctx, SRemoteValueNode* pRes, SSDataBlock* pBlock) {
22,511,321✔
4150
  int32_t code = 0;
22,511,321✔
4151
  bool needFree = true;
22,511,321✔
4152
  int32_t colNum = taosArrayGetSize(pBlock->pDataBlock);
22,511,150✔
4153
  if (NULL == pBlock->pDataBlock || 1 != colNum || pBlock->info.rows > 1) {
22,510,941✔
4154
    qError("%s invalid scl fetch res block, pDataBlock:%p, colNum:%d, rows:%" PRId64, 
171✔
4155
      ctx->idStr, pBlock->pDataBlock, colNum, pBlock->info.rows);
4156
    return TSDB_CODE_PAR_INVALID_SCALAR_SUBQ_RES_ROWS;
×
4157
  }
4158
  
4159
  pRes->val.node.type = QUERY_NODE_VALUE;
22,510,171✔
4160
  pRes->val.flag &= (~VALUE_FLAG_VAL_UNSET);
22,509,729✔
4161
  pRes->val.translate = true;
22,510,770✔
4162
  
4163
  SColumnInfoData* pCol = taosArrayGet(pBlock->pDataBlock, 0);
22,510,437✔
4164
  if (colDataIsNull_s(pCol, 0)) {
22,512,205✔
4165
    pRes->val.isNull = true;
1,819,567✔
4166
  } else {
4167
    code = nodesSetValueNodeValueExt(&pRes->val, colDataGetData(pCol, 0), &needFree);
20,692,638✔
4168
  }
4169

4170
  if (!needFree) {
22,509,441✔
4171
    pCol->pData = NULL;
17,164✔
4172
  }
4173

4174
  return code;
22,509,441✔
4175
}
4176

4177
int32_t remoteFetchCallBack(void* param, SDataBuf* pMsg, int32_t code) {
42,805,781✔
4178
  SScalarFetchParam* pParam = (SScalarFetchParam*)param;
42,805,781✔
4179
  STaskSubJobCtx* ctx = pParam->pSubJobCtx;
42,805,781✔
4180
  SSDataBlock* pResBlock = NULL;
42,806,872✔
4181
  
4182
  taosMemoryFreeClear(pMsg->pEpSet);
42,807,421✔
4183

4184
  if (NULL == ctx) {
42,805,090✔
4185
    qWarn("scl fetch ctx not exists since it may have been released");
1,198✔
4186
    goto _exit;
1,198✔
4187
  }
4188

4189
  qDebug("%s subQIdx %d got rsp, code:%d, rsp:%p", ctx->idStr, pParam->subQIdx, code, pMsg->pData);
42,803,892✔
4190

4191
  taosWLockLatch(&ctx->lock);
42,803,892✔
4192
  ctx->param = NULL;
42,806,603✔
4193
  taosWUnLockLatch(&ctx->lock);
42,806,603✔
4194

4195
  if (ctx->transporterId > 0) {
42,807,756✔
4196
    int32_t ret = asyncFreeConnById(ctx->rpcHandle, ctx->transporterId);
42,805,579✔
4197
    if (ret != 0) {
42,806,777✔
4198
      qDebug("%s failed to free subQ rpc handle, code:%s, subQIdx:%d", ctx->idStr, tstrerror(ret), pParam->subQIdx);
×
4199
    }
4200
    ctx->transporterId = -1;
42,806,777✔
4201
  }
4202

4203
  if (0 == code && NULL == pMsg->pData) {
42,807,801✔
4204
    qError("%s invalid rsp msg, msgType:%d, len:%d", ctx->idStr, pMsg->msgType, pMsg->len);
×
4205
    code = TSDB_CODE_QRY_INVALID_MSG;
×
4206
  }
4207

4208
  if (code == TSDB_CODE_SUCCESS) {
42,804,488✔
4209
    SRetrieveTableRsp* pRsp = pMsg->pData;
34,720,937✔
4210
    pRsp->numOfRows = htobe64(pRsp->numOfRows);
34,720,937✔
4211
    pRsp->compLen = htonl(pRsp->compLen);
34,721,474✔
4212
    pRsp->payloadLen = htonl(pRsp->payloadLen);
34,720,844✔
4213
    pRsp->numOfCols = htonl(pRsp->numOfCols);
34,722,453✔
4214
    pRsp->useconds = htobe64(pRsp->useconds);
34,721,916✔
4215
    pRsp->numOfBlocks = htonl(pRsp->numOfBlocks);
34,719,865✔
4216

4217
    if (pRsp->numOfRows > 1 || pRsp->numOfBlocks > 1 || !pRsp->completed) {
34,720,212✔
4218
      qError("%s invalid scl fetch rsp received, subQIdx:%d, rows:%" PRId64 ", blocks:%d, completed:%d", 
762,825✔
4219
        ctx->idStr, pParam->subQIdx, pRsp->numOfRows, pRsp->numOfBlocks, pRsp->completed);
4220
      ctx->code = TSDB_CODE_PAR_INVALID_SCALAR_SUBQ_RES_ROWS;
763,267✔
4221
    } else if (0 == pRsp->numOfRows) {
33,957,541✔
4222
      SRemoteValueNode* pRemote = (SRemoteValueNode*)pParam->pRes;
11,444,371✔
4223
      pRemote->val.node.type = QUERY_NODE_VALUE;
11,445,445✔
4224
      pRemote->val.isNull = true;
11,444,908✔
4225
      pRemote->val.translate = true;
11,445,445✔
4226
      pRemote->val.flag &= (~VALUE_FLAG_VAL_UNSET);
11,445,445✔
4227
      taosArraySet(ctx->subResValues, pParam->subQIdx, &pParam->pRes);
11,445,445✔
4228
    } else {
4229
      qDebug("%s scl fetch rsp received, subQIdx:%d, rows:%" PRId64 , ctx->idStr, pParam->subQIdx, pRsp->numOfRows);
22,511,195✔
4230
      ctx->code = createBlockFromRemoteValueNode(&pResBlock, pParam->pRes);
22,511,195✔
4231
      if (TSDB_CODE_SUCCESS == ctx->code) {
22,511,713✔
4232
        ctx->code = blockDataEnsureCapacity(pResBlock, 1);
22,511,808✔
4233
      }
4234
      if (TSDB_CODE_SUCCESS == ctx->code) {
22,513,246✔
4235
        ctx->code = extractSingleRspBlock(pRsp, pResBlock);
22,512,709✔
4236
      }
4237
      if (TSDB_CODE_SUCCESS == ctx->code) {
22,511,687✔
4238
        ctx->code = setValueFromResBlock(ctx, pParam->pRes, pResBlock);
22,509,287✔
4239
      }
4240
      if (TSDB_CODE_SUCCESS == ctx->code) {
22,513,072✔
4241
        taosArraySet(ctx->subResValues, pParam->subQIdx, &pParam->pRes);
22,509,441✔
4242
      }
4243
    }
4244
  } else {
4245
    ctx->code = rpcCvtErrCode(code);
8,083,551✔
4246
    if (ctx->code != code) {
8,085,303✔
4247
      qError("%s scl fetch rsp received, subQIdx:%d, error:%s, cvted error: %s", ctx->idStr, pParam->subQIdx,
×
4248
             tstrerror(code), tstrerror(ctx->code));
4249
    } else {
4250
      qError("%s scl fetch rsp received, subQIdx:%d, error:%s", ctx->idStr, pParam->subQIdx, tstrerror(code));
8,083,551✔
4251
    }
4252
  }
4253
  
4254
  code = tsem_post(&pParam->pSubJobCtx->ready);
42,805,132✔
4255
  if (code != TSDB_CODE_SUCCESS) {
42,807,756✔
4256
    qError("failed to invoke post when scl fetch rsp is ready, code:%s", tstrerror(code));
×
4257
  }
4258

4259
_exit:
42,808,954✔
4260

4261
  taosMemoryFree(pMsg->pData);
42,806,478✔
4262
  blockDataDestroy(pResBlock);
42,808,070✔
4263

4264
  return code;
42,806,112✔
4265
}
4266

4267

4268
int32_t fetchRemoteValueImpl(STaskSubJobCtx* ctx, int32_t subQIdx, SRemoteValueNode* pRes) {
42,796,830✔
4269
  int32_t          code = TSDB_CODE_SUCCESS;
42,796,830✔
4270
  int32_t          lino = 0;
42,796,830✔
4271
  SDownstreamSourceNode* pSource = (SDownstreamSourceNode*)taosArrayGetP(ctx->subEndPoints, subQIdx);
42,796,830✔
4272

4273
  SResFetchReq req = {0};
42,800,040✔
4274
  req.header.vgId = pSource->addr.nodeId;
42,800,040✔
4275
  req.sId = pSource->sId;
42,797,580✔
4276
  req.clientId = pSource->clientId;
42,797,942✔
4277
  req.taskId = pSource->taskId;
42,800,499✔
4278
  req.queryId = ctx->queryId;
42,794,134✔
4279
  req.execId = pSource->execId;
42,783,788✔
4280

4281
  int32_t msgSize = tSerializeSResFetchReq(NULL, 0, &req, false);
42,791,902✔
4282
  if (msgSize < 0) {
42,788,049✔
4283
    return msgSize;
×
4284
  }
4285

4286
  void* msg = taosMemoryCalloc(1, msgSize);
42,788,049✔
4287
  if (NULL == msg) {
42,779,218✔
4288
    return terrno;
×
4289
  }
4290

4291
  msgSize = tSerializeSResFetchReq(msg, msgSize, &req, false);
42,779,218✔
4292
  if (msgSize < 0) {
42,792,026✔
4293
    taosMemoryFree(msg);
×
4294
    return msgSize;
×
4295
  }
4296

4297
  qDebug("%s scl build fetch msg and send to nodeId:%d, ep:%s, clientId:0x%" PRIx64 " taskId:0x%" PRIx64
42,792,026✔
4298
         ", execId:%d",
4299
         ctx->idStr, pSource->addr.nodeId, pSource->addr.epSet.eps[0].fqdn, pSource->clientId,
4300
         pSource->taskId, pSource->execId);
4301

4302
  // send the fetch remote task result reques
4303
  SMsgSendInfo* pMsgSendInfo = taosMemoryCalloc(1, sizeof(SMsgSendInfo));
42,787,198✔
4304
  if (NULL == pMsgSendInfo) {
42,778,937✔
4305
    taosMemoryFreeClear(msg);
×
4306
    qError("%s prepare message %d failed", ctx->idStr, (int32_t)sizeof(SMsgSendInfo));
×
4307
    return terrno;
×
4308
  }
4309

4310
  SScalarFetchParam* param = taosMemoryMalloc(sizeof(SScalarFetchParam));
42,778,937✔
4311
  if (NULL == param) {
42,784,904✔
4312
    taosMemoryFreeClear(msg);
×
4313
    taosMemoryFreeClear(pMsgSendInfo);
×
4314
    qError("%s prepare param %d failed", ctx->idStr, (int32_t)sizeof(SScalarFetchParam));
×
4315
    return terrno;
×
4316
  }
4317

4318
  taosWLockLatch(&ctx->lock);
42,784,904✔
4319
  
4320
  if (ctx->code) {
42,803,130✔
4321
    qError("task has been killed, error:%s", tstrerror(ctx->code));
×
4322
    taosMemoryFree(param);
×
4323
    code = ctx->code;
×
4324
    goto _end;
×
4325
  } else {
4326
    ctx->param = param;
42,788,496✔
4327
  }
4328
  
4329
  taosWUnLockLatch(&ctx->lock);
42,794,335✔
4330

4331
  param->subQIdx = subQIdx;
42,800,048✔
4332
  param->pRes = pRes;
42,800,205✔
4333
  param->pSubJobCtx = ctx;
42,797,855✔
4334

4335
  pMsgSendInfo->param = param;
42,789,962✔
4336
  pMsgSendInfo->paramFreeFp = taosAutoMemoryFree;
42,783,295✔
4337
  pMsgSendInfo->msgInfo.pData = msg;
42,784,658✔
4338
  pMsgSendInfo->msgInfo.len = msgSize;
42,776,412✔
4339
  pMsgSendInfo->msgType = pSource->fetchMsgType;
42,793,436✔
4340
  pMsgSendInfo->fp = remoteFetchCallBack;
42,778,624✔
4341
  pMsgSendInfo->requestId = ctx->queryId;
42,775,499✔
4342

4343
  code = asyncSendMsgToServer(ctx->rpcHandle, &pSource->addr.epSet, &ctx->transporterId, pMsgSendInfo);
42,784,561✔
4344
  QUERY_CHECK_CODE(code, lino, _end);
42,806,996✔
4345

4346
  code = qSemWait(ctx->pTaskInfo, &ctx->ready);
42,806,996✔
4347
  if (isTaskKilled(ctx->pTaskInfo)) {
42,808,954✔
4348
    code = getTaskCode(ctx->pTaskInfo);
4,193✔
4349
  } else {
4350
    code = ctx->code;
42,804,162✔
4351
  }
4352
      
4353
_end:
42,808,954✔
4354

4355
  taosWLockLatch(&ctx->lock);
42,808,954✔
4356
  ctx->param = NULL;
42,808,355✔
4357
  taosWUnLockLatch(&ctx->lock);
42,808,355✔
4358

4359
  if (code != TSDB_CODE_SUCCESS) {
42,808,954✔
4360
    qError("%s %s failed at line %d since %s", ctx->idStr, __func__, lino, tstrerror(code));
8,850,263✔
4361
  }
4362
  return code;
42,808,954✔
4363
}
4364

4365

4366
int32_t qFetchRemoteValue(void* pCtx, int32_t subQIdx, SRemoteValueNode* pRes) {
44,227,095✔
4367
  STaskSubJobCtx*  ctx = (STaskSubJobCtx*)pCtx;
44,227,095✔
4368
  int32_t code = 0, lino = 0;
44,227,095✔
4369
  int32_t       subEndPoinsNum = taosArrayGetSize(ctx->subEndPoints);
44,227,095✔
4370
  if (subQIdx >= subEndPoinsNum) {
44,227,453✔
4371
    qError("%s invalid subQIdx %d, subEndPointsNum:%d", ctx->idStr, subQIdx, subEndPoinsNum);
×
4372
    return TSDB_CODE_QRY_SUBQ_NOT_FOUND;
×
4373
  }
4374

4375
  SValueNode** ppRes = taosArrayGet(ctx->subResValues, subQIdx);
44,227,453✔
4376
  if (NULL == *ppRes) {
44,231,329✔
4377
    TAOS_CHECK_EXIT(fetchRemoteValueImpl(ctx, subQIdx, pRes));
42,802,687✔
4378
    *ppRes = (SValueNode*)pRes;
33,958,691✔
4379
  } else {
4380
    TAOS_CHECK_EXIT(valueNodeCopy(*ppRes, &pRes->val));
1,431,339✔
4381
    pRes->val.node.type = QUERY_NODE_VALUE;
1,431,339✔
4382
  }
4383

4384
_exit:
44,240,293✔
4385

4386
  if (code) {
44,240,293✔
4387
    qError("%s %s failed at line %d since %s", ctx->idStr, __func__, lino, tstrerror(code));
8,850,263✔
4388
  }
4389

4390
  return code;
44,239,238✔
4391
}
4392

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