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

taosdata / TDengine / #5071

17 May 2026 01:15AM UTC coverage: 63.054% (-10.3%) from 73.326%
#5071

push

travis-ci

web-flow
feat (TDgpt): Dynamic Model Synchronization Enhancements (#35344)

* refactor: do some internal refactor.

* fix: fix multiprocess sync issue.

* feat: add dynamic anomaly detection and forecasting services

* fix: log error message for undeploying model in exception handling

* Potential fix for pull request finding

Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com>

* Potential fix for pull request finding

Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com>

* Potential fix for pull request finding

Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com>

* Potential fix for pull request finding

Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com>

* fix: handle undeploy when model exists only on disk

Agent-Logs-Url: https://github.com/taosdata/TDengine/sessions/286aafa0-c3ce-4c27-b803-2707571e9dc1

Co-authored-by: hjxilinx <8252296+hjxilinx@users.noreply.github.com>

* fix: guard dynamic registry concurrent access

Agent-Logs-Url: https://github.com/taosdata/TDengine/sessions/5e4db858-6458-40f4-ac28-d1b1b7f97c18

Co-authored-by: hjxilinx <8252296+hjxilinx@users.noreply.github.com>

* fix: tighten service list locking scope

Agent-Logs-Url: https://github.com/taosdata/TDengine/sessions/5e4db858-6458-40f4-ac28-d1b1b7f97c18

Co-authored-by: hjxilinx <8252296+hjxilinx@users.noreply.github.com>

* fix: restore prophet support and update tests per review feedback

Agent-Logs-Url: https://github.com/taosdata/TDengine/sessions/92298ae1-7da6-4d07-b20e-101c7cd0b26b

Co-authored-by: hjxilinx <8252296+hjxilinx@users.noreply.github.com>

* fix: improve test name and move copy inside lock scope

Agent-Logs-Url: https://github.com/taosdata/TDengine/sessions/92298ae1-7da6-4d07-b20e-101c7cd0b26b

Co-authored-by: hjxilinx <8252296+hjxilinx@users.noreply.github.com>

* Potential fix for pull request finding

Co-au... (continued)

238317 of 377957 relevant lines covered (63.05%)

130539817.12 hits per line

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

82.95
/source/libs/executor/src/countwindowoperator.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 "executorInt.h"
17
#include "filter.h"
18
#include "function.h"
19
#include "functionMgt.h"
20
#include "operator.h"
21
#include "querytask.h"
22
#include "tcommon.h"
23
#include "tcompare.h"
24
#include "tdatablock.h"
25
#include "ttime.h"
26

27
typedef struct SCountWindowResult {
28
  int32_t    winRows;
29
  TSKEY      winSKey;
30
  SResultRow row;
31
} SCountWindowResult;
32

33
typedef struct SCountWindowSupp {
34
  SArray* pWinStates;
35
  int32_t stateIndex;
36
  int32_t curStateIndex;
37
  TSKEY   lastTs; // this ts is used to record the last timestamp, so that we can know whether the new row's ts is duplicated
38
} SCountWindowSupp;
39

40
typedef struct SCountWindowOperatorInfo {
41
  SOptrBasicInfo     binfo;
42
  SAggSupporter      aggSup;
43
  SExprSupp          scalarSup;
44
  int32_t            tsSlotId;  // primary timestamp column slot id
45
  STimeWindowAggSupp twAggSup;
46
  uint64_t           groupId;  // current group id, used to identify the data block from different groups
47
  SResultRow*        pRow;
48
  int32_t            windowCount;
49
  int32_t            windowSliding;
50
  SCountWindowSupp   countSup;
51
  SSDataBlock*       pPreDataBlock;
52
  int32_t            preStateIndex;
53
  bool               indefRowsMode;
54
  SIndefRowsRuntime  indefRows;
55
  struct SOperatorInfo* pOperator;
56
} SCountWindowOperatorInfo;
57

58
void destroyCountWindowOperatorInfo(void* param) {
292,579✔
59
  SCountWindowOperatorInfo* pInfo = (SCountWindowOperatorInfo*)param;
292,579✔
60
  if (pInfo == NULL) {
292,579✔
61
    return;
×
62
  }
63
  cleanupBasicInfo(&pInfo->binfo);
292,579✔
64
  colDataDestroy(&pInfo->twAggSup.timeWindowData);
292,579✔
65
  cleanupIndefRowsRuntime(&pInfo->indefRows, pInfo->pOperator);
292,579✔
66

67
  cleanupAggSup(&pInfo->aggSup);
292,579✔
68
  cleanupExprSupp(&pInfo->scalarSup);
292,579✔
69
  taosArrayDestroy(pInfo->countSup.pWinStates);
292,579✔
70
  taosMemoryFreeClear(param);
292,579✔
71
}
72

73
static int32_t countWindowAggregateNext(SOperatorInfo* pOperator, SSDataBlock** ppRes);
74

75
static void clearWinStateBuff(SCountWindowResult* pBuff) {
137,766,178✔
76
  pBuff->winRows = 0;
137,766,178✔
77
  pBuff->winSKey = 0;
137,766,626✔
78
}
137,766,626✔
79

80
static SCountWindowResult* getCountWinStateInfo(SCountWindowSupp* pCountSup) {
140,671,184✔
81
  SCountWindowResult* pBuffInfo = taosArrayGet(pCountSup->pWinStates, pCountSup->stateIndex);
140,671,184✔
82
  pCountSup->curStateIndex = pCountSup->stateIndex;
140,671,184✔
83
  if (!pBuffInfo) {
140,671,184✔
84
    terrno = TSDB_CODE_QRY_EXECUTOR_INTERNAL_ERROR;
×
85
    qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
86
    return NULL;
×
87
  }
88
  int32_t size = taosArrayGetSize(pCountSup->pWinStates);
140,671,184✔
89
  if (size == 0) {
140,671,184✔
90
    terrno = TSDB_CODE_QRY_EXECUTOR_INTERNAL_ERROR;
×
91
    qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
92
    return NULL;
×
93
  }
94
  pCountSup->stateIndex = (pCountSup->stateIndex + 1) % size;
140,671,184✔
95
  return pBuffInfo;
140,671,184✔
96
}
97

98
static int32_t setCountWindowOutputBuff(SExprSupp* pExprSup, SCountWindowSupp* pCountSup, SResultRow** pResult,
140,585,265✔
99
                                        SCountWindowResult** ppResBuff) {
100
  int32_t             code = TSDB_CODE_SUCCESS;
140,585,265✔
101
  int32_t             lino = 0;
140,585,265✔
102
  SCountWindowResult* pBuff = getCountWinStateInfo(pCountSup);
140,585,265✔
103
  QUERY_CHECK_NULL(pBuff, code, lino, _end, terrno);
140,585,265✔
104
  (*pResult) = &pBuff->row;
140,585,265✔
105
  code = setResultRowInitCtx(*pResult, pExprSup->pCtx, pExprSup->numOfExprs, pExprSup->rowEntryInfoOffset);
140,585,265✔
106
  (*ppResBuff) = pBuff;
140,585,713✔
107

108
_end:
140,585,713✔
109
  if (code != TSDB_CODE_SUCCESS) {
140,585,713✔
110
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
111
  }
112
  return code;
140,585,713✔
113
}
114

115
static int32_t updateCountWindowInfo(int32_t start, int32_t blockRows, int32_t countWinRows, int32_t* pCurrentRows) {
140,112,229✔
116
  int32_t rows = TMIN(countWinRows - (*pCurrentRows), blockRows - start);
140,112,229✔
117
  (*pCurrentRows) += rows;
140,112,229✔
118
  return rows;
140,112,229✔
119
}
120

121
void doCountWindowAggImpl(SOperatorInfo* pOperator, SSDataBlock* pBlock) {
3,144,537✔
122
  int32_t                   code = TSDB_CODE_SUCCESS;
3,144,537✔
123
  int32_t                   lino = 0;
3,144,537✔
124
  SExecTaskInfo*            pTaskInfo = pOperator->pTaskInfo;
3,144,537✔
125
  SExprSupp*                pExprSup = &pOperator->exprSupp;
3,144,537✔
126
  SCountWindowOperatorInfo* pInfo = pOperator->info;
3,144,537✔
127
  SSDataBlock*              pRes = pInfo->binfo.pRes;
3,144,537✔
128
  SColumnInfoData*          pColInfoData = taosArrayGet(pBlock->pDataBlock, pInfo->tsSlotId);
3,144,537✔
129
  QUERY_CHECK_NULL(pColInfoData, code, lino, _end, terrno);
3,144,537✔
130
  TSKEY* tsCols = (TSKEY*)pColInfoData->pData;
3,144,537✔
131
  int32_t numOfBuff = taosArrayGetSize(pInfo->countSup.pWinStates);
3,144,537✔
132
  if (numOfBuff == 0) {
3,144,537✔
133
    code = TSDB_CODE_QRY_EXECUTOR_INTERNAL_ERROR;
×
134
    qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(code));
×
135
    T_LONG_JMP(pTaskInfo->env, code);
×
136
  }
137
  pInfo->countSup.stateIndex = (pInfo->preStateIndex + 1) % numOfBuff;
3,144,537✔
138

139
  int32_t newSize = pRes->info.rows + pBlock->info.rows / pInfo->windowSliding + 1;
3,144,537✔
140
  if (!pInfo->indefRowsMode && newSize > pRes->info.capacity) {
3,144,537✔
141
    code = blockDataEnsureCapacity(pRes, newSize);
×
142
    QUERY_CHECK_CODE(code, lino, _end);
×
143
  }
144

145
  for (int32_t i = 0; i < pBlock->info.rows; i++) {
2,147,483,647✔
146
    if (pBlock->info.scanFlag != PRE_SCAN) {
2,147,483,647✔
147
      if (pInfo->countSup.lastTs == INT64_MIN) {
2,147,483,647✔
148
        pInfo->countSup.lastTs = tsCols[i];
358,173✔
149
      } else {
150
        if (tsCols[i] == pInfo->countSup.lastTs) {
2,147,483,647✔
151
          qError("duplicate timestamp found in count window operator" PRId64 ", timestamp: %" PRId64, tsCols[i]);
11,834✔
152
          code = TSDB_CODE_QRY_WINDOW_DUP_TIMESTAMP;
6,278✔
153
          QUERY_CHECK_CODE(code, lino, _end);
6,278✔
154
        } else {
155
          pInfo->countSup.lastTs = tsCols[i];
2,147,483,647✔
156
        }
157
      }
158
    }
159
  }
160

161
  for (int32_t i = 0; i < pBlock->info.rows;) {
143,244,036✔
162
    SCountWindowResult* pBuffInfo = NULL;
140,111,781✔
163
    if (!pInfo->indefRowsMode) {
140,111,781✔
164
      code = setCountWindowOutputBuff(pExprSup, &pInfo->countSup, &pInfo->pRow, &pBuffInfo);
140,052,812✔
165
      if (code != TSDB_CODE_SUCCESS) {
140,053,260✔
166
        qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
167
        T_LONG_JMP(pTaskInfo->env, code);
×
168
      }
169
    } else {
170
      pBuffInfo = getCountWinStateInfo(&pInfo->countSup);
58,969✔
171
      QUERY_CHECK_NULL(pBuffInfo, code, lino, _end, terrno);
58,969✔
172
    }
173
    int32_t prevRows = pBuffInfo->winRows;
140,112,229✔
174
    int32_t num = updateCountWindowInfo(i, pBlock->info.rows, pInfo->windowCount, &pBuffInfo->winRows);
140,112,229✔
175
    int32_t step = num;
140,112,229✔
176
    if (prevRows == 0) {
140,112,229✔
177
      pBuffInfo->winSKey = tsCols[i];
137,766,626✔
178
      if (!pInfo->indefRowsMode) {
137,766,626✔
179
        pInfo->pRow->win.skey = tsCols[i];
137,707,657✔
180
      }
181
    }
182
    STimeWindow win = {.skey = pBuffInfo->winSKey, .ekey = tsCols[num + i - 1]};
140,112,229✔
183
    if (pInfo->indefRowsMode) {
140,112,229✔
184
      SIndefRowsWindowState* pState = NULL;
58,969✔
185
      code = applyIndefRowsFuncOnWindowState(pOperator, &pInfo->indefRows, &pState, pInfo->binfo.pRes,
58,969✔
186
                                             pBlock->info.id.groupId, &win, pBlock, i, num, pInfo->binfo.inputTsOrder,
187
                                             pInfo->aggSup.resultRowSize);
188
      QUERY_CHECK_CODE(code, lino, _end);
58,969✔
189
    } else {
190
      pInfo->pRow->win.ekey = win.ekey;
140,052,812✔
191
      updateTimeWindowInfo(&pInfo->twAggSup.timeWindowData, &pInfo->pRow->win, 0);
140,052,812✔
192
      code = applyAggFunctionOnPartialTuples(pTaskInfo, pExprSup->pCtx, &pInfo->twAggSup.timeWindowData, i, num,
280,106,072✔
193
                                             pBlock->info.rows, pExprSup->numOfExprs);
140,052,812✔
194
      QUERY_CHECK_CODE(code, lino, _end);
140,053,260✔
195
    }
196
    if (pInfo->windowCount != pInfo->windowSliding) {
140,112,229✔
197
      if (prevRows <= pInfo->windowSliding) {
50,724✔
198
        if (pBuffInfo->winRows > pInfo->windowSliding) {
51,172✔
199
          step = pInfo->windowSliding - prevRows;
41,627✔
200
        } else {
201
          step = pInfo->windowSliding;
9,097✔
202
        }
203
      } else {
204
        step = 0;
×
205
      }
206
    }
207
    if (pBuffInfo->winRows == pInfo->windowCount) {
140,111,781✔
208
      if (pInfo->indefRowsMode) {
137,471,838✔
209
        SIndefRowsWindowState* pState = findIndefRowsWindowState(&pInfo->indefRows, pBlock->info.id.groupId,
42,752✔
210
                                                                 pBuffInfo->winSKey);
42,752✔
211
        QUERY_CHECK_NULL(pState, code, lino, _end, TSDB_CODE_QRY_WINDOW_STATE_NOT_EXIST);
42,752✔
212
        code = closeIndefRowsWindowState(pOperator, &pInfo->indefRows, pState);
42,752✔
213
        QUERY_CHECK_CODE(code, lino, _end);
42,752✔
214
      } else {
215
        doUpdateNumOfRows(pExprSup->pCtx, pInfo->pRow, pExprSup->numOfExprs, pExprSup->rowEntryInfoOffset);
137,428,638✔
216
        code = copyResultrowToDataBlock(pExprSup->pExprInfo, pExprSup->numOfExprs, pInfo->pRow, pExprSup->pCtx, pRes,
137,428,638✔
217
                                        pExprSup->rowEntryInfoOffset, pTaskInfo);
137,429,086✔
218
        QUERY_CHECK_CODE(code, lino, _end);
137,429,086✔
219
        pRes->info.rows += pInfo->pRow->numOfRows;
137,429,086✔
220
      }
221
      clearWinStateBuff(pBuffInfo);
137,471,390✔
222
      pInfo->preStateIndex = pInfo->countSup.curStateIndex;
137,471,390✔
223
      if (!pInfo->indefRowsMode) {
137,471,390✔
224
        clearResultRowInitFlag(pExprSup->pCtx, pExprSup->numOfExprs);
137,428,638✔
225
      }
226
    }
227
    i += step;
140,111,333✔
228
  }
229

230
  if (!pInfo->indefRowsMode) {
3,132,703✔
231
    code = doFilter(pRes, pOperator->exprSupp.pFilterInfo, NULL, NULL);
3,111,327✔
232
    QUERY_CHECK_CODE(code, lino, _end);
3,111,327✔
233
  }
234

235
_end:
3,132,703✔
236
  if (code != TSDB_CODE_SUCCESS) {
3,144,537✔
237
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
11,834✔
238
    pTaskInfo->code = code;
11,834✔
239
    T_LONG_JMP(pTaskInfo->env, code);
11,834✔
240
  }
241
}
3,132,703✔
242

243
static void buildCountResult(SOperatorInfo* pOperator, SCountWindowOperatorInfo* pInfo, uint64_t groupId, SSDataBlock* pBlock) {
497,927✔
244
  SExprSupp*        pExprSup = &pOperator->exprSupp;
497,927✔
245
  SCountWindowSupp* pCountSup = &pInfo->countSup;
497,927✔
246
  SExecTaskInfo*    pTaskInfo = pOperator->pTaskInfo;
497,927✔
247
  SFilterInfo*      pFilterInfo = pOperator->exprSupp.pFilterInfo;
497,927✔
248
  SResultRow* pResultRow = NULL;
497,927✔
249
  int32_t     code = TSDB_CODE_SUCCESS;
497,927✔
250
  int32_t     lino = 0;
497,927✔
251
  int32_t     numOfBuff = taosArrayGetSize(pCountSup->pWinStates);
497,927✔
252
  int32_t     newSize = pBlock->info.rows + numOfBuff;
497,927✔
253
  if (!pInfo->indefRowsMode && newSize > pBlock->info.capacity) {
497,927✔
254
    code = blockDataEnsureCapacity(pBlock, newSize);
×
255
    QUERY_CHECK_CODE(code, lino, _end);
×
256
  }
257
  pCountSup->stateIndex = (pInfo->preStateIndex + 1) % numOfBuff;
497,927✔
258
  for (int32_t i = 0; i < numOfBuff; i++) {
1,057,330✔
259
    SCountWindowResult* pBuff = NULL;
559,403✔
260
    if (!pInfo->indefRowsMode) {
559,403✔
261
      code = setCountWindowOutputBuff(pExprSup, pCountSup, &pResultRow, &pBuff);
532,453✔
262
      QUERY_CHECK_CODE(code, lino, _end);
532,453✔
263
    } else {
264
      pBuff = getCountWinStateInfo(pCountSup);
26,950✔
265
      QUERY_CHECK_NULL(pBuff, code, lino, _end, terrno);
26,950✔
266
    }
267
    if (pBuff->winRows == 0) {
559,403✔
268
      continue;
264,615✔
269
    }
270
    if (pInfo->indefRowsMode) {
294,788✔
271
      SIndefRowsWindowState* pState = findIndefRowsWindowState(&pInfo->indefRows, groupId, pBuff->winSKey);
16,217✔
272
      QUERY_CHECK_NULL(pState, code, lino, _end, TSDB_CODE_QRY_WINDOW_STATE_NOT_EXIST);
16,217✔
273
      code = closeIndefRowsWindowState(pOperator, &pInfo->indefRows, pState);
16,217✔
274
      QUERY_CHECK_CODE(code, lino, _end);
16,217✔
275
    } else {
276
      doUpdateNumOfRows(pExprSup->pCtx, pResultRow, pExprSup->numOfExprs, pExprSup->rowEntryInfoOffset);
278,571✔
277
      code = copyResultrowToDataBlock(pExprSup->pExprInfo, pExprSup->numOfExprs, pResultRow, pExprSup->pCtx, pBlock,
278,571✔
278
                                      pExprSup->rowEntryInfoOffset, pTaskInfo);
278,571✔
279
      QUERY_CHECK_CODE(code, lino, _end);
278,571✔
280
      pBlock->info.rows += pResultRow->numOfRows;
278,571✔
281
    }
282
    clearWinStateBuff(pBuff);
294,788✔
283
    if (!pInfo->indefRowsMode) {
294,788✔
284
      clearResultRowInitFlag(pExprSup->pCtx, pExprSup->numOfExprs);
278,571✔
285
    }
286
  }
287
  if (!pInfo->indefRowsMode) {
497,927✔
288
    code = doFilter(pBlock, pFilterInfo, NULL, NULL);
470,977✔
289
    QUERY_CHECK_CODE(code, lino, _end);
470,977✔
290
  }
291

292
_end:
497,927✔
293
  if (code != TSDB_CODE_SUCCESS) {
497,927✔
294
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
295
    T_LONG_JMP(pTaskInfo->env, code);
×
296
  }
297
}
497,927✔
298

299
static int32_t countWindowAggregateNext(SOperatorInfo* pOperator, SSDataBlock** ppRes) {
569,263✔
300
  int32_t                   code = TSDB_CODE_SUCCESS;
569,263✔
301
  int32_t                   lino = 0;
569,263✔
302
  SCountWindowOperatorInfo* pInfo = pOperator->info;
569,263✔
303
  SExecTaskInfo*            pTaskInfo = pOperator->pTaskInfo;
569,263✔
304
  SExprSupp*                pExprSup = &pOperator->exprSupp;
568,703✔
305
  int32_t                   order = pInfo->binfo.inputTsOrder;
568,703✔
306
  SSDataBlock*              pRes = pInfo->binfo.pRes;
568,703✔
307

308
  if (pOperator->status == OP_EXEC_DONE) {
569,263✔
309
    *ppRes = NULL;
×
310
    return code;
×
311
  }
312

313
  if (pInfo->indefRowsMode) {
569,263✔
314
    (*ppRes) = getNextIndefRowsResultBlock(&pInfo->indefRows, pOperator);
71,091✔
315
    if ((*ppRes) != NULL) {
71,091✔
316
      return code;
21,845✔
317
    }
318
  }
319

320
  blockDataCleanup(pRes);
547,418✔
321

322
  while (1) {
3,111,327✔
323
    SSDataBlock* pBlock = NULL;
3,658,745✔
324
    if (pInfo->pPreDataBlock == NULL) { 
3,658,745✔
325
      pBlock = getNextBlockFromDownstream(pOperator, 0);
3,532,017✔
326
    } else {
327
      pBlock = pInfo->pPreDataBlock;
126,728✔
328
      pInfo->pPreDataBlock = NULL;
126,728✔
329
    }
330

331
    if (pBlock == NULL) {
3,658,745✔
332
      break;
365,725✔
333
    }
334

335
    pRes->info.scanFlag = pBlock->info.scanFlag;
3,293,020✔
336
    code = setInputDataBlock(pExprSup, pBlock, order, MAIN_SCAN, true);
3,293,020✔
337
    QUERY_CHECK_CODE(code, lino, _end);
3,293,020✔
338

339
    code = blockDataUpdateTsWindow(pBlock, pInfo->tsSlotId);
3,293,020✔
340
    QUERY_CHECK_CODE(code, lino, _end);
3,293,020✔
341

342
    // there is an scalar expression that needs to be calculated right before apply the group aggregation.
343
    if (pInfo->scalarSup.pExprInfo != NULL) {
3,293,020✔
344
      code = projectApplyFunctions(pInfo->scalarSup.pExprInfo, pBlock, pBlock, pInfo->scalarSup.pCtx,
118,332✔
345
                                   pInfo->scalarSup.numOfExprs, NULL, GET_STM_RTINFO(pOperator->pTaskInfo), pOperator->pTaskInfo);
59,166✔
346
      QUERY_CHECK_CODE(code, lino, _end);
59,166✔
347
    }
348

349
    if (pInfo->groupId == 0) {
3,276,739✔
350
      pInfo->groupId = pBlock->info.id.groupId;
2,251,449✔
351
    } else if (pInfo->groupId != pBlock->info.id.groupId) {
1,025,290✔
352
      pInfo->pPreDataBlock = pBlock;
132,202✔
353
      pRes->info.id.groupId = pInfo->groupId;
132,202✔
354
      buildCountResult(pOperator, pInfo, pInfo->groupId, pRes);
132,202✔
355
      pInfo->groupId = pBlock->info.id.groupId;
132,202✔
356
      pInfo->countSup.lastTs = INT64_MIN;
132,202✔
357
      if (pInfo->indefRowsMode) {
132,202✔
358
        (*ppRes) = getNextIndefRowsResultBlock(&pInfo->indefRows, pOperator);
11,112✔
359
        if ((*ppRes) != NULL) {
11,112✔
360
          return code;
11,112✔
361
        }
362
        continue;
×
363
      }
364
      if (pRes->info.rows > 0) {
121,090✔
365
        (*ppRes) = pRes;
121,090✔
366
        return code;
121,090✔
367
      }
368
    }
369

370
    doCountWindowAggImpl(pOperator, pBlock);
3,144,537✔
371
    if (pInfo->indefRowsMode) {
3,132,703✔
372
      (*ppRes) = getNextIndefRowsResultBlock(&pInfo->indefRows, pOperator);
21,376✔
373
      if ((*ppRes) != NULL) {
21,376✔
374
        return code;
20,907✔
375
      }
376
      continue;
469✔
377
    }
378

379
    if (pRes->info.rows >= pOperator->resultInfo.threshold) {
3,111,327✔
380
      pRes->info.id.groupId = pInfo->groupId;
469✔
381
      (*ppRes) = pRes;
469✔
382
      return code;
469✔
383
    }
384
  }
385

386
  pRes->info.id.groupId = pInfo->groupId;
365,725✔
387
  buildCountResult(pOperator, pInfo, pInfo->groupId, pRes);
365,725✔
388

389
  if (pInfo->indefRowsMode) {
365,725✔
390
    (*ppRes) = getNextIndefRowsResultBlock(&pInfo->indefRows, pOperator);
15,838✔
391
    if ((*ppRes) == NULL) {
15,838✔
392
      setOperatorCompleted(pOperator);
10,733✔
393
    }
394
    return code;
15,838✔
395
  }
396

397
_end:
349,887✔
398
  if (code != TSDB_CODE_SUCCESS) {
366,168✔
399
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
16,281✔
400
    pTaskInfo->code = code;
16,281✔
401
    T_LONG_JMP(pTaskInfo->env, code);
16,281✔
402
  }
403
  (*ppRes) = pRes->info.rows == 0 ? NULL : pRes;
349,887✔
404
  return code;
349,887✔
405
}
406

407
static int32_t resetCountWindowOperatorState(SOperatorInfo* pOper) {
×
408
  SCountWindowOperatorInfo* pCount = pOper->info;
×
409
  SExecTaskInfo*           pTaskInfo = pOper->pTaskInfo;
×
410
  SCountWindowPhysiNode* pPhynode = (SCountWindowPhysiNode*)pOper->pPhyNode;
×
411
  pOper->status = OP_NOT_OPENED;
×
412
  
413
  resetBasicOperatorState(&pCount->binfo);
×
414
  pCount->groupId = 0;
×
415
  pCount->countSup.stateIndex = 0;
×
416
  pCount->countSup.lastTs = INT64_MIN;
×
417
  pCount->pPreDataBlock = NULL;
×
418
  pCount->preStateIndex = 0;
×
419
  pCount->pRow = NULL;
×
420
  resetIndefRowsRuntime(&pCount->indefRows, pOper);
×
421
  int32_t numOfBuff = taosArrayGetSize(pCount->countSup.pWinStates);
×
422
  for (int32_t i = 0; i < numOfBuff; ++i) {
×
423
    SCountWindowResult* pBuff = taosArrayGet(pCount->countSup.pWinStates, i);
×
424
    if (pBuff != NULL) {
×
425
      clearWinStateBuff(pBuff);
×
426
      memset(&pBuff->row, 0, pCount->aggSup.resultRowSize);
×
427
    }
428
  }
429

430
  colDataDestroy(&pCount->twAggSup.timeWindowData);
×
431
  int32_t code = initExecTimeWindowInfo(&pCount->twAggSup.timeWindowData, &pTaskInfo->window);
×
432
  
433
  if (code == 0) {
×
434
    code = resetAggSup(&pOper->exprSupp, &pCount->aggSup, pTaskInfo, pPhynode->window.pFuncs, NULL,
×
435
                        sizeof(int64_t) * 2 + POINTER_BYTES, pTaskInfo->id.str, NULL,
×
436
                        &pTaskInfo->storageAPI.functionStore);
437
  }
438
  if (code == 0) {
×
439
    code = resetExprSupp(&pCount->scalarSup, pTaskInfo, pPhynode->window.pExprs, NULL,
×
440
                          &pTaskInfo->storageAPI.functionStore);
441
  }
442
  return code;
×
443
}
444

445
int32_t createCountwindowOperatorInfo(SOperatorInfo* downstream, SPhysiNode* physiNode,
291,966✔
446
                                             SExecTaskInfo* pTaskInfo, SOperatorInfo** pOptrInfo) {
447
  QRY_PARAM_CHECK(pOptrInfo);
291,966✔
448

449
  int32_t                   code = TSDB_CODE_SUCCESS;
292,579✔
450
  int32_t                   lino = 0;
292,579✔
451
  SCountWindowOperatorInfo* pInfo = taosMemoryCalloc(1, sizeof(SCountWindowOperatorInfo));
292,579✔
452
  SOperatorInfo*            pOperator = taosMemoryCalloc(1, sizeof(SOperatorInfo));
292,579✔
453
  if (pInfo == NULL || pOperator == NULL) {
292,579✔
454
    code = terrno;
×
455
    goto _error;
×
456
  }
457
  initOperatorCostInfo(pOperator);
292,579✔
458

459
  pOperator->pPhyNode = physiNode;
292,579✔
460
  pOperator->exprSupp.hasWindowOrGroup = true;
292,579✔
461
  pOperator->exprSupp.hasWindow = true;
292,579✔
462

463
  SCountWindowPhysiNode* pCountWindowNode = (SCountWindowPhysiNode*)physiNode;
292,579✔
464

465
  pInfo->tsSlotId = ((SColumnNode*)pCountWindowNode->window.pTspk)->slotId;
292,579✔
466

467
  if (pCountWindowNode->window.pExprs != NULL) {
292,579✔
468
    int32_t    numOfScalarExpr = 0;
69,889✔
469
    SExprInfo* pScalarExprInfo = NULL;
69,889✔
470
    code = createExprInfo(pCountWindowNode->window.pExprs, NULL, &pScalarExprInfo, &numOfScalarExpr);
69,889✔
471
    QUERY_CHECK_CODE(code, lino, _error);
69,889✔
472
    code = initExprSupp(&pInfo->scalarSup, pScalarExprInfo, numOfScalarExpr, &pTaskInfo->storageAPI.functionStore);
69,889✔
473
    QUERY_CHECK_CODE(code, lino, _error);
69,889✔
474
  }
475

476
  size_t     keyBufSize = 0;
292,579✔
477
  int32_t    num = 0;
292,579✔
478
  SExprInfo* pExprInfo = NULL;
292,579✔
479
  code = createExprInfo(pCountWindowNode->window.pFuncs, NULL, &pExprInfo, &num);
292,579✔
480
  QUERY_CHECK_CODE(code, lino, _error);
292,579✔
481

482
  initResultSizeInfo(&pOperator->resultInfo, 4096);
292,579✔
483

484
  code = initAggSup(&pOperator->exprSupp, &pInfo->aggSup, pExprInfo, num, keyBufSize, pTaskInfo->id.str,
292,579✔
485
                    NULL, &pTaskInfo->storageAPI.functionStore);
486
  QUERY_CHECK_CODE(code, lino, _error);
292,579✔
487

488
  pInfo->indefRowsMode = pCountWindowNode->window.indefRowsFunc;
292,579✔
489
  if (pInfo->indefRowsMode) {
292,579✔
490
    code = initIndefRowsRuntime(&pInfo->indefRows, pOperator->exprSupp.pCtx, num, pOperator->resultInfo.capacity);
12,122✔
491
    QUERY_CHECK_CODE(code, lino, _error);
12,122✔
492
  }
493

494
  SSDataBlock* pResBlock = createDataBlockFromDescNode(pCountWindowNode->window.node.pOutputDataBlockDesc);
292,579✔
495
  QUERY_CHECK_NULL(pResBlock, code, lino, _error, terrno);
292,579✔
496
  initBasicInfo(&pInfo->binfo, pResBlock);
292,579✔
497

498
  code = blockDataEnsureCapacity(pResBlock, pOperator->resultInfo.capacity);
292,579✔
499
  QUERY_CHECK_CODE(code, lino, _error);
292,579✔
500

501
  initResultRowInfo(&pInfo->binfo.resultRowInfo);
292,579✔
502
  pInfo->binfo.inputTsOrder = physiNode->inputTsOrder;
292,579✔
503
  pInfo->binfo.outputTsOrder = physiNode->outputTsOrder;
292,579✔
504
  pInfo->windowCount = pCountWindowNode->windowCount;
292,579✔
505
  pInfo->windowSliding = pCountWindowNode->windowSliding;
292,579✔
506
  // sizeof(SCountWindowResult)
507
  int32_t itemSize = sizeof(SCountWindowResult) - sizeof(SResultRow) + pInfo->aggSup.resultRowSize;
292,579✔
508
  int32_t numOfItem = 1;
292,579✔
509
  if (pInfo->windowCount != pInfo->windowSliding) {
292,579✔
510
    numOfItem = pInfo->windowCount / pInfo->windowSliding + 1;
19,401✔
511
  }
512

513
  pInfo->countSup.pWinStates = taosArrayInit_s(itemSize, numOfItem);
292,579✔
514
  if (!pInfo->countSup.pWinStates) {
292,579✔
515
    goto _error;
×
516
  }
517

518
  pInfo->countSup.stateIndex = 0;
292,579✔
519
  pInfo->pPreDataBlock = NULL;
292,579✔
520
  pInfo->preStateIndex = 0;
292,579✔
521
  pInfo->countSup.lastTs = INT64_MIN;
292,579✔
522
  pInfo->pOperator = pOperator;
292,579✔
523

524
  code = filterInitFromNode((SNode*)pCountWindowNode->window.node.pConditions, &pOperator->exprSupp.pFilterInfo, 0,
292,579✔
525
                            pTaskInfo->pStreamRuntimeInfo);
292,579✔
526
  QUERY_CHECK_CODE(code, lino, _error);
292,579✔
527
  filterSetExecContext(pOperator->exprSupp.pFilterInfo, pTaskInfo, isTaskKilled);
292,579✔
528

529
  code = initExecTimeWindowInfo(&pInfo->twAggSup.timeWindowData, &pTaskInfo->window);
292,579✔
530
  QUERY_CHECK_CODE(code, lino, _error);
292,579✔
531

532
  setOperatorInfo(pOperator, "CountWindowOperator", QUERY_NODE_PHYSICAL_PLAN_MERGE_COUNT, true, OP_NOT_OPENED, pInfo,
292,579✔
533
                  pTaskInfo);
534
  pOperator->fpSet = createOperatorFpSet(optrDummyOpenFn, countWindowAggregateNext, NULL, destroyCountWindowOperatorInfo,
292,579✔
535
                                         optrDefaultBufFn, NULL, optrDefaultGetNextExtFn, NULL);
536
  setOperatorResetStateFn(pOperator, resetCountWindowOperatorState);                                    
292,579✔
537
  code = appendDownstream(pOperator, &downstream, 1);
292,579✔
538
  if (code != TSDB_CODE_SUCCESS) {
292,579✔
539
    goto _error;
×
540
  }
541

542
  *pOptrInfo = pOperator;
292,579✔
543
  return TSDB_CODE_SUCCESS;
292,579✔
544

545
_error:
×
546
  if (pInfo != NULL) {
×
547
    destroyCountWindowOperatorInfo(pInfo);
×
548
  }
549

550
  destroyOperatorAndDownstreams(pOperator, &downstream, 1);
×
551
  pTaskInfo->code = code;
×
552
  return code;
×
553
}
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