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

taosdata / TDengine / #5080

17 May 2026 01:15AM UTC coverage: 73.403% (+0.008%) from 73.395%
#5080

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)

281716 of 383795 relevant lines covered (73.4%)

139531599.76 hits per line

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

73.46
/source/dnode/vnode/src/tq/tqMeta.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
#include "tdbInt.h"
16
#include "tq.h"
17

18
int32_t tqEncodeSTqHandle(SEncoder* pEncoder, const STqHandle* pHandle) {
2,808,309✔
19
  if (pEncoder == NULL || pHandle == NULL) {
2,808,309✔
20
    return TSDB_CODE_INVALID_PARA;
×
21
  }
22
  int32_t code = 0;
2,810,138✔
23
  int32_t lino;
24

25
  TAOS_CHECK_EXIT(tStartEncode(pEncoder));
2,810,138✔
26
  TAOS_CHECK_EXIT(tEncodeCStr(pEncoder, pHandle->subKey));
5,618,367✔
27
  TAOS_CHECK_EXIT(tEncodeI8(pEncoder, pHandle->fetchMeta));
5,617,876✔
28
  TAOS_CHECK_EXIT(tEncodeI64(pEncoder, pHandle->consumerId));
5,617,901✔
29
  TAOS_CHECK_EXIT(tEncodeI64(pEncoder, pHandle->snapshotVer));
5,614,999✔
30
  TAOS_CHECK_EXIT(tEncodeI32(pEncoder, pHandle->epoch));
5,610,878✔
31
  TAOS_CHECK_EXIT(tEncodeI8(pEncoder, pHandle->execHandle.subType));
5,612,745✔
32
  if (pHandle->execHandle.subType == TOPIC_SUB_TYPE__COLUMN) {
2,807,504✔
33
    TAOS_CHECK_EXIT(tEncodeCStr(pEncoder, pHandle->execHandle.execCol.qmsg));
4,555,375✔
34
    TAOS_CHECK_EXIT(tEncodeSSchemaWrapper(pEncoder, &pHandle->execHandle.execCol.pSW));
4,556,926✔
35
  } else if (pHandle->execHandle.subType == TOPIC_SUB_TYPE__DB) {
529,858✔
36
    TAOS_CHECK_EXIT(tEncodeI32(pEncoder, 0));
434,986✔
37
  } else if (pHandle->execHandle.subType == TOPIC_SUB_TYPE__TABLE) {
95,284✔
38
    TAOS_CHECK_EXIT(tEncodeI64(pEncoder, pHandle->execHandle.execTb.suid));
190,568✔
39
    if (pHandle->execHandle.execTb.qmsg != NULL) {
95,672✔
40
      TAOS_CHECK_EXIT(tEncodeCStr(pEncoder, pHandle->execHandle.execTb.qmsg));
190,956✔
41
    }
42
  }
43
  tEndEncode(pEncoder);
2,809,514✔
44
_exit:
2,801,110✔
45
  if (code) {
2,801,110✔
46
    return code;
×
47
  } else {
48
    return pEncoder->pos;
2,801,110✔
49
  }
50
}
51

52
int32_t tqDecodeSTqHandle(SDecoder* pDecoder, STqHandle* pHandle) {
81,222✔
53
  if (pDecoder == NULL || pHandle == NULL) {
81,222✔
54
    return TSDB_CODE_INVALID_PARA;
×
55
  }
56
  int32_t code = 0;
81,222✔
57
  int32_t lino;
58

59
  TAOS_CHECK_EXIT(tStartDecode(pDecoder));
81,222✔
60
  TAOS_CHECK_EXIT(tDecodeCStrTo(pDecoder, pHandle->subKey));
81,222✔
61
  TAOS_CHECK_EXIT(tDecodeI8(pDecoder, &pHandle->fetchMeta));
162,444✔
62
  TAOS_CHECK_EXIT(tDecodeI64(pDecoder, &pHandle->consumerId));
162,444✔
63
  TAOS_CHECK_EXIT(tDecodeI64(pDecoder, &pHandle->snapshotVer));
162,444✔
64
  TAOS_CHECK_EXIT(tDecodeI32(pDecoder, &pHandle->epoch));
162,444✔
65
  TAOS_CHECK_EXIT(tDecodeI8(pDecoder, &pHandle->execHandle.subType));
162,444✔
66
  if (pHandle->execHandle.subType == TOPIC_SUB_TYPE__COLUMN) {
81,222✔
67
    TAOS_CHECK_EXIT(tDecodeCStrAlloc(pDecoder, &pHandle->execHandle.execCol.qmsg));
109,032✔
68
    if (!tDecodeIsEnd(pDecoder)) {
54,516✔
69
      TAOS_CHECK_EXIT(tDecodeSSchemaWrapper(pDecoder, &pHandle->execHandle.execCol.pSW));
109,032✔
70
    }
71
  } else if (pHandle->execHandle.subType == TOPIC_SUB_TYPE__DB) {
26,706✔
72
    int32_t size = 0;
18,362✔
73
    TAOS_CHECK_EXIT(tDecodeI32(pDecoder, &size));
18,362✔
74
    for (int32_t i = 0; i < size; i++) {
18,362✔
75
      int64_t tbUid = 0;
×
76
      TAOS_CHECK_EXIT(tDecodeI64(pDecoder, &tbUid));
×
77
    }
78
  } else if (pHandle->execHandle.subType == TOPIC_SUB_TYPE__TABLE) {
8,344✔
79
    TAOS_CHECK_EXIT(tDecodeI64(pDecoder, &pHandle->execHandle.execTb.suid));
16,688✔
80
    if (!tDecodeIsEnd(pDecoder)) {
8,344✔
81
      TAOS_CHECK_EXIT(tDecodeCStrAlloc(pDecoder, &pHandle->execHandle.execTb.qmsg));
16,688✔
82
    }
83
  }
84
  tEndDecode(pDecoder);
81,222✔
85

86
_exit:
81,222✔
87
  return code;
81,222✔
88
}
89

90
void tqDestroySTqHandle(void* data) {
551,232✔
91
  if (data == NULL) return;
551,232✔
92
  STqHandle* pData = (STqHandle*)data;
551,232✔
93
  qDestroyTask(pData->execHandle.task);
551,232✔
94

95
  if (pData->execHandle.subType == TOPIC_SUB_TYPE__COLUMN) {
549,458✔
96
    taosMemoryFreeClear(pData->execHandle.execCol.qmsg);
397,153✔
97
    taosMemoryFreeClear(pData->execHandle.execCol.pSW.pSchema);
397,312✔
98
  } else if (pData->execHandle.subType == TOPIC_SUB_TYPE__DB) {
151,973✔
99
    tqReaderClose(pData->execHandle.pTqReader);
121,457✔
100
    walCloseReader(pData->pWalReader);
121,457✔
101
  } else if (pData->execHandle.subType == TOPIC_SUB_TYPE__TABLE) {
30,838✔
102
    walCloseReader(pData->pWalReader);
30,838✔
103
    tqReaderClose(pData->execHandle.pTqReader);
30,838✔
104
    taosMemoryFreeClear(pData->execHandle.execTb.qmsg);
30,838✔
105
    nodesDestroyNode(pData->execHandle.execTb.node);
30,838✔
106
  }
107
  if (pData->msg != NULL) {
548,181✔
108
    rpcFreeCont(pData->msg->pCont);
×
109
    taosMemoryFree(pData->msg);
×
110
    pData->msg = NULL;
×
111
  }
112
  if (pData->block != NULL) {
550,135✔
113
    blockDataDestroy(pData->block);
×
114
  }
115
  if (pData->pRef) {
551,232✔
116
    walCloseRef(pData->pRef->pWal, pData->pRef->refId);
535,669✔
117
  }
118
  taosHashCleanup(pData->tableCreateTimeHash);
549,310✔
119
}
120

121
int32_t tqMetaDecodeOffsetInfo(STqOffset* info, void* pVal, uint32_t vLen) {
79,229✔
122
  if (info == NULL || pVal == NULL) {
79,229✔
123
    return TSDB_CODE_INVALID_PARA;
×
124
  }
125
  SDecoder decoder = {0};
79,229✔
126
  tDecoderInit(&decoder, (uint8_t*)pVal, vLen);
79,229✔
127
  int32_t code = tDecodeSTqOffset(&decoder, info);
79,229✔
128
  tDecoderClear(&decoder);
79,229✔
129

130
  if (code != 0) {
79,229✔
131
    tDeleteSTqOffset(info);
×
132
    return TSDB_CODE_OUT_OF_MEMORY;
×
133
  }
134
  return code;
79,229✔
135
}
136

137
int32_t tqMetaSaveOffset(STQ* pTq, STqOffset* pOffset) {
274,418✔
138
  if (pTq == NULL || pOffset == NULL) {
274,418✔
139
    return TSDB_CODE_INVALID_PARA;
×
140
  }
141
  void*    buf = NULL;
274,418✔
142
  int32_t  code = TDB_CODE_SUCCESS;
274,418✔
143
  uint32_t  vlen;
144
  SEncoder encoder = {0};
274,418✔
145
  tEncodeSize(tEncodeSTqOffset, pOffset, vlen, code);
274,418✔
146
  if (code < 0) {
274,418✔
147
    goto END;
×
148
  }
149

150
  buf = taosMemoryCalloc(1, vlen);
274,418✔
151
  if (buf == NULL) {
274,418✔
152
    code = terrno;
×
153
    goto END;
×
154
  }
155

156
  tEncoderInit(&encoder, buf, vlen);
274,418✔
157
  code = tEncodeSTqOffset(&encoder, pOffset);
274,418✔
158
  if (code < 0) {
274,418✔
159
    goto END;
×
160
  }
161

162
  taosWLockLatch(&pTq->lock);
274,418✔
163
  code = tqMetaSaveInfo(pTq, pTq->pOffsetStore, pOffset->subKey, strlen(pOffset->subKey), buf, vlen);
274,418✔
164
  taosWUnLockLatch(&pTq->lock);
274,418✔
165
END:
274,418✔
166
  tEncoderClear(&encoder);
274,418✔
167
  taosMemoryFree(buf);
274,418✔
168
  return code;
274,418✔
169
}
170

171
int32_t tqMetaSaveInfo(STQ* pTq, TTB* ttb, const void* key, uint32_t kLen, const void* value, uint32_t vLen) {
1,705,142✔
172
  if (pTq == NULL || ttb == NULL || key == NULL || value == NULL) {
1,705,142✔
173
    return TSDB_CODE_INVALID_PARA;
×
174
  }
175
  int32_t code = TDB_CODE_SUCCESS;
1,707,843✔
176
  TXN*    txn = NULL;
1,707,843✔
177

178
  TQ_ERR_GO_TO_END(
1,707,843✔
179
      tdbBegin(pTq->pMetaDB, &txn, tdbDefaultMalloc, tdbDefaultFree, NULL, TDB_TXN_WRITE | TDB_TXN_READ_UNCOMMITTED));
180
  TQ_ERR_GO_TO_END(tdbTbUpsert(ttb, key, kLen, value, vLen, txn));
1,704,093✔
181
  TQ_ERR_GO_TO_END(tdbCommit(pTq->pMetaDB, txn));
1,704,099✔
182
  TQ_ERR_GO_TO_END(tdbPostCommit(pTq->pMetaDB, txn));
1,703,205✔
183

184
  return 0;
1,702,659✔
185

186
END:
×
187
  if (txn != NULL) {
×
188
    tdbAbort(pTq->pMetaDB, txn);
×
189
  }
190
  return code;
×
191
}
192

193
int32_t tqMetaDeleteInfo(STQ* pTq, TTB* ttb, const void* key, uint32_t kLen) {
684,208✔
194
  if (pTq == NULL || ttb == NULL || key == NULL) {
684,208✔
195
    return TSDB_CODE_INVALID_PARA;
×
196
  }
197
  int32_t code = TDB_CODE_SUCCESS;
684,560✔
198
  TXN*    txn = NULL;
684,560✔
199

200
  TQ_ERR_GO_TO_END(
684,560✔
201
      tdbBegin(pTq->pMetaDB, &txn, tdbDefaultMalloc, tdbDefaultFree, NULL, TDB_TXN_WRITE | TDB_TXN_READ_UNCOMMITTED));
202
  TQ_ERR_GO_TO_END(tdbTbDelete(ttb, key, kLen, txn));
680,808✔
203
  TQ_ERR_GO_TO_END(tdbCommit(pTq->pMetaDB, txn));
360,389✔
204
  TQ_ERR_GO_TO_END(tdbPostCommit(pTq->pMetaDB, txn));
359,207✔
205

206
  return 0;
361,832✔
207

208
END:
317,056✔
209
  if (txn != NULL) {
317,812✔
210
    tdbAbort(pTq->pMetaDB, txn);
318,934✔
211
  }
212
  return code;
319,651✔
213
}
214

215
int32_t tqMetaGetOffset(STQ* pTq, const char* subkey, STqOffset** pOffset) {
17,183,913✔
216
  if (pTq == NULL || subkey == NULL || pOffset == NULL) {
17,183,913✔
217
    return TSDB_CODE_INVALID_PARA;
713✔
218
  }
219
  void* data = taosHashGet(pTq->pOffset, subkey, strlen(subkey));
17,183,200✔
220
  if (data == NULL) {
17,182,478✔
221
    int vLen = 0;
5,893,858✔
222
    taosRLockLatch(&pTq->lock);
5,893,858✔
223
    if (tdbTbGet(pTq->pOffsetStore, subkey, strlen(subkey), &data, &vLen) < 0) {
5,895,654✔
224
      taosRUnLockLatch(&pTq->lock);
5,809,344✔
225
      tdbFree(data);
5,818,359✔
226
      return TSDB_CODE_OUT_OF_MEMORY;
5,789,481✔
227
    }
228
    taosRUnLockLatch(&pTq->lock);
65,047✔
229

230
    STqOffset offset = {0};
65,047✔
231
    if (tqMetaDecodeOffsetInfo(&offset, data, vLen >= 0 ? vLen : 0) != TDB_CODE_SUCCESS) {
65,047✔
232
      tdbFree(data);
×
233
      return TSDB_CODE_OUT_OF_MEMORY;
×
234
    }
235

236
    if (taosHashPut(pTq->pOffset, subkey, strlen(subkey), &offset, sizeof(STqOffset)) != 0) {
65,047✔
237
      tDeleteSTqOffset(&offset);
×
238
      tdbFree(data);
×
239
      return terrno;
×
240
    }
241
    tdbFree(data);
65,047✔
242

243
    *pOffset = taosHashGet(pTq->pOffset, subkey, strlen(subkey));
65,047✔
244
    if (*pOffset == NULL) {
65,047✔
245
      return TSDB_CODE_OUT_OF_MEMORY;
×
246
    }
247
  } else {
248
    *pOffset = data;
11,288,620✔
249
  }
250
  return 0;
11,353,667✔
251
}
252

253
int32_t tqMetaSaveHandle(STQ* pTq, const char* key, const STqHandle* pHandle) {
1,404,877✔
254
  if (pTq == NULL || key == NULL || pHandle == NULL) {
1,404,877✔
255
    return TSDB_CODE_INVALID_PARA;
×
256
  }
257
  int32_t  code = TDB_CODE_SUCCESS;
1,404,877✔
258
  uint32_t  vlen;
259
  void*    buf = NULL;
1,404,877✔
260
  SEncoder encoder = {0};
1,404,877✔
261
  tEncodeSize(tqEncodeSTqHandle, pHandle, vlen, code);
1,404,877✔
262
  if (code < 0) {
1,404,496✔
263
    goto END;
×
264
  }
265

266
  tqDebug("tq save %s(%d) handle consumer:0x%" PRIx64 " epoch:%d vgId:%d", pHandle->subKey,
1,404,496✔
267
          (int32_t)strlen(pHandle->subKey), pHandle->consumerId, pHandle->epoch, TD_VID(pTq->pVnode));
268

269
  buf = taosMemoryCalloc(1, vlen);
1,405,261✔
270
  if (buf == NULL) {
1,405,261✔
271
    code = terrno;
×
272
    goto END;
×
273
  }
274

275
  tEncoderInit(&encoder, buf, vlen);
1,405,261✔
276
  code = tqEncodeSTqHandle(&encoder, pHandle);
1,405,261✔
277
  if (code < 0) {
1,404,848✔
278
    goto END;
×
279
  }
280

281
  TQ_ERR_GO_TO_END(tqMetaSaveInfo(pTq, pTq->pExecStore, key, strlen(key), buf, vlen));
1,404,848✔
282

283
END:
1,400,136✔
284
  tEncoderClear(&encoder);
1,400,536✔
285
  taosMemoryFree(buf);
1,397,099✔
286
  return code;
1,392,306✔
287
}
288

289
static int tqMetaInitHandle(STQ* pTq, STqHandle* handle) {
537,150✔
290
  if (pTq == NULL || handle == NULL) {
537,150✔
291
    return TSDB_CODE_INVALID_PARA;
×
292
  }
293
  int32_t code = TDB_CODE_SUCCESS;
537,150✔
294
  SArray* tbUidList = NULL;
537,150✔
295

296
  SVnode* pVnode = pTq->pVnode;
536,789✔
297
  int32_t vgId = TD_VID(pVnode);
537,150✔
298

299
  handle->pRef = walOpenRef(pVnode->pWal);
536,789✔
300
  TQ_NULL_GO_TO_END(handle->pRef);
537,150✔
301

302
  SReadHandle reader = {0};
536,775✔
303
  reader = (SReadHandle){
536,414✔
304
      .vnode = pVnode,
305
      .initTableReader = true,
306
      .initTqReader = true,
307
      .version = handle->snapshotVer,
536,775✔
308
  };
309

310
  initStorageAPI(&reader.api);
536,414✔
311

312
  if (handle->execHandle.subType == TOPIC_SUB_TYPE__COLUMN) {
535,953✔
313
    handle->execHandle.task = qCreateQueueExecTaskInfo(handle->execHandle.execCol.qmsg, &reader, vgId, handle->consumerId);
389,987✔
314
    TQ_NULL_GO_TO_END(handle->execHandle.task);
390,375✔
315
    void* scanner = NULL;
390,375✔
316
    qExtractTmqScanner(handle->execHandle.task, &scanner);
390,375✔
317
    TQ_NULL_GO_TO_END(scanner);
390,375✔
318
    handle->execHandle.pTqReader = qExtractReaderFromTmqScanner(scanner);
390,375✔
319
    TQ_NULL_GO_TO_END(handle->execHandle.pTqReader);
390,375✔
320
  } else if (handle->execHandle.subType == TOPIC_SUB_TYPE__DB) {
146,775✔
321
    handle->pWalReader = walOpenReader(pVnode->pWal, 0);
118,697✔
322
    TQ_NULL_GO_TO_END(handle->pWalReader);
118,697✔
323
    handle->execHandle.pTqReader = tqReaderOpen(pVnode);
118,697✔
324
    TQ_NULL_GO_TO_END(handle->execHandle.pTqReader);
118,375✔
325
    TQ_ERR_GO_TO_END(buildSnapContext(reader.vnode, reader.version, 0, handle->execHandle.subType, handle->fetchMeta,
118,375✔
326
                                      (SSnapContext**)(&reader.sContext)));
327
    handle->execHandle.task = qCreateQueueExecTaskInfo(NULL, &reader, vgId, handle->consumerId);
118,697✔
328
    TQ_NULL_GO_TO_END(handle->execHandle.task);
118,697✔
329
  } else if (handle->execHandle.subType == TOPIC_SUB_TYPE__TABLE) {
28,078✔
330
    handle->pWalReader = walOpenReader(pVnode->pWal, 0);
28,078✔
331
    TQ_NULL_GO_TO_END(handle->pWalReader);
28,078✔
332
    if (handle->execHandle.execTb.qmsg != NULL && strcmp(handle->execHandle.execTb.qmsg, "") != 0) {
28,078✔
333
      if (nodesStringToNode(handle->execHandle.execTb.qmsg, &handle->execHandle.execTb.node) != 0) {
10,780✔
334
        tqError("nodesStringToNode error in sub stable, since %s", terrstr());
×
335
        return TSDB_CODE_SCH_INTERNAL_ERROR;
×
336
      }
337
    }
338
    TQ_ERR_GO_TO_END(buildSnapContext(reader.vnode, reader.version, handle->execHandle.execTb.suid,
28,078✔
339
                                      handle->execHandle.subType, handle->fetchMeta,
340
                                      (SSnapContext**)(&reader.sContext)));
341
    handle->execHandle.task = qCreateQueueExecTaskInfo(NULL, &reader, vgId, handle->consumerId);
28,078✔
342
    TQ_NULL_GO_TO_END(handle->execHandle.task);
28,078✔
343
    TQ_ERR_GO_TO_END(qGetTableList(handle->execHandle.execTb.suid, pVnode, handle->execHandle.execTb.node, &tbUidList,
27,690✔
344
                                handle->execHandle.task));
345
    tqInfo("vgId:%d, tq try to get ctb for stb subscribe, suid:%" PRId64, pVnode->config.vgId,
28,078✔
346
           handle->execHandle.execTb.suid);
347
    handle->execHandle.pTqReader = tqReaderOpen(pVnode);
28,078✔
348
    TQ_NULL_GO_TO_END(handle->execHandle.pTqReader);
28,078✔
349
    TQ_ERR_GO_TO_END(tqReaderSetTbUidList(handle->execHandle.pTqReader, tbUidList, NULL));
28,078✔
350
  }
351
  handle->tableCreateTimeHash = (SHashObj*)taosHashInit(64, taosGetDefaultHashFunction(TSDB_DATA_TYPE_BIGINT), true, HASH_ENTRY_LOCK);
537,150✔
352

353
END:
536,814✔
354
  reader.api.snapshotFn.destroySnapshot(reader.sContext);
536,401✔
355
  taosArrayDestroy(tbUidList);
536,737✔
356
  return code;
536,737✔
357
}
358

359
static int32_t tqMetaRestoreHandle(STQ* pTq, void* pVal, uint32_t vLen, STqHandle* handle) {
67,140✔
360
  if (pTq == NULL || pVal == NULL || handle == NULL) {
67,140✔
361
    return TSDB_CODE_INVALID_PARA;
×
362
  }
363
  int32_t  vgId = TD_VID(pTq->pVnode);
67,140✔
364
  SDecoder decoder = {0};
67,140✔
365
  int32_t  code = TDB_CODE_SUCCESS;
67,140✔
366

367
  tDecoderInit(&decoder, (uint8_t*)pVal, vLen);
67,140✔
368
  TQ_ERR_GO_TO_END(tqDecodeSTqHandle(&decoder, handle));
67,140✔
369
  TQ_ERR_GO_TO_END(tqMetaInitHandle(pTq, handle));
67,140✔
370
  tqInfo("tqMetaRestoreHandle %s consumer 0x%" PRIx64 " vgId:%d", handle->subKey, handle->consumerId, vgId);
67,140✔
371
  code = taosHashPut(pTq->pHandle, handle->subKey, strlen(handle->subKey), handle, sizeof(STqHandle));
67,140✔
372

373
END:
67,140✔
374
  tDecoderClear(&decoder);
67,140✔
375
  return code;
67,140✔
376
}
377

378
int32_t tqMetaCreateHandle(STQ* pTq, SMqRebVgReq* req, STqHandle* handle) {
470,010✔
379
  if (pTq == NULL || req == NULL || handle == NULL) {
470,010✔
380
    return TSDB_CODE_INVALID_PARA;
×
381
  }
382
  int32_t vgId = TD_VID(pTq->pVnode);
470,010✔
383

384
  (void)memcpy(handle->subKey, req->subKey, TSDB_SUBSCRIBE_KEY_LEN);
470,010✔
385
  handle->consumerId = req->newConsumerId;
470,010✔
386

387
  handle->execHandle.subType = req->subType;
470,010✔
388
  handle->fetchMeta = req->withMeta;
470,010✔
389
  if (req->subType == TOPIC_SUB_TYPE__COLUMN) {
469,622✔
390
    void* tmp = taosStrdup(req->qmsg);
344,033✔
391
    if (tmp == NULL) {
344,421✔
392
      return terrno;
×
393
    }
394
    handle->execHandle.execCol.qmsg = tmp;
344,421✔
395
    handle->execHandle.execCol.pSW = req->schema;
344,421✔
396
    req->schema.pSchema = NULL;
344,421✔
397
  } else if (req->subType == TOPIC_SUB_TYPE__TABLE) {
125,589✔
398
    handle->execHandle.execTb.suid = req->suid;
22,494✔
399
    void* tmp = taosStrdup(req->qmsg);
22,494✔
400
    if (tmp == NULL) {
22,494✔
401
      return terrno;
×
402
    }
403
    handle->execHandle.execTb.qmsg = tmp;
22,494✔
404
  }
405

406
  handle->snapshotVer = walGetCommittedVer(pTq->pVnode->pWal);
470,010✔
407

408
  int32_t code = tqMetaInitHandle(pTq, handle);
469,649✔
409
  if (code != 0) {
470,010✔
410
    return code;
×
411
  }
412
  tqInfo("tqMetaCreateHandle %s consumer 0x%" PRIx64 " vgId:%d, snapshotVer:%" PRId64, handle->subKey,
470,010✔
413
         handle->consumerId, vgId, handle->snapshotVer);
414
  return taosHashPut(pTq->pHandle, handle->subKey, strlen(handle->subKey), handle, sizeof(STqHandle));
470,423✔
415
}
416

417
static int32_t tqMetaTransformInfo(TDB* pMetaDB, TTB* pOld, TTB* pNew) {
×
418
  if (pMetaDB == NULL || pOld == NULL || pNew == NULL) {
×
419
    return TSDB_CODE_INVALID_PARA;
×
420
  }
421
  TBC*  pCur = NULL;
×
422
  void* pKey = NULL;
×
423
  int   kLen = 0;
×
424
  void* pVal = NULL;
×
425
  int   vLen = 0;
×
426
  TXN*  txn = NULL;
×
427

428
  int32_t code = TDB_CODE_SUCCESS;
×
429

430
  TQ_ERR_GO_TO_END(tdbTbcOpen(pOld, &pCur, NULL));
×
431
  TQ_ERR_GO_TO_END(
×
432
      tdbBegin(pMetaDB, &txn, tdbDefaultMalloc, tdbDefaultFree, NULL, TDB_TXN_WRITE | TDB_TXN_READ_UNCOMMITTED));
433

434
  TQ_ERR_GO_TO_END(tdbTbcMoveToFirst(pCur));
×
435
  while (tdbTbcNext(pCur, &pKey, &kLen, &pVal, &vLen) == 0) {
×
436
    TQ_ERR_GO_TO_END(tdbTbUpsert(pNew, pKey, kLen, pVal, vLen, txn));
×
437
  }
438

439
  TQ_ERR_GO_TO_END(tdbCommit(pMetaDB, txn));
×
440
  TQ_ERR_GO_TO_END(tdbPostCommit(pMetaDB, txn));
×
441

442
END:
×
443
  tdbFree(pKey);
×
444
  tdbFree(pVal);
×
445
  tdbTbcClose(pCur);
×
446
  return code;
×
447
}
448

449
int32_t tqMetaGetHandle(STQ* pTq, const char* key, STqHandle** pHandle) {
322,594,706✔
450
  if (pTq == NULL || key == NULL || pHandle == NULL) {
322,594,706✔
451
    return TSDB_CODE_INVALID_PARA;
×
452
  }
453
  void* data = taosHashGet(pTq->pHandle, key, strlen(key));
322,683,084✔
454
  if (data == NULL) {
322,594,031✔
455
    int vLen = 0;
537,745✔
456
    if (tdbTbGet(pTq->pExecStore, key, (int)strlen(key), &data, &vLen) < 0) {
537,745✔
457
      tdbFree(data);
469,182✔
458
      return TSDB_CODE_MND_SUBSCRIBE_NOT_EXIST;
469,895✔
459
    }
460
    STqHandle handle = {0};
67,140✔
461
    int32_t code = tqMetaRestoreHandle(pTq, data, vLen >= 0 ? vLen : 0, &handle);
67,140✔
462
    if (code != 0) {
67,140✔
463
      tdbFree(data);
×
464
      tqDestroySTqHandle(&handle);
×
465
      return code;
×
466
    }
467
    tdbFree(data);
67,140✔
468
    *pHandle = taosHashGet(pTq->pHandle, key, strlen(key));
67,140✔
469
    if (*pHandle == NULL) {
67,140✔
470
      return terrno;
×
471
    }
472
  } else {
473
    *pHandle = data;
322,056,286✔
474
  }
475
  return TDB_CODE_SUCCESS;
322,199,949✔
476
}
477

478
int32_t tqMetaOpenTdb(STQ* pTq) {
4,888,914✔
479
  if (pTq == NULL) {
4,888,914✔
480
    return TSDB_CODE_INVALID_PARA;
×
481
  }
482
  int32_t code = TDB_CODE_SUCCESS;
4,888,914✔
483
  TQ_ERR_GO_TO_END(tdbOpen(pTq->path, 16 * 1024, 1, &pTq->pMetaDB, 0, NULL));
4,888,914✔
484
  TQ_ERR_GO_TO_END(tdbTbOpen("tq.db", -1, -1, NULL, pTq->pMetaDB, &pTq->pExecStore, 0));
4,892,219✔
485
  TQ_ERR_GO_TO_END(tdbTbOpen("tq.offset.db", -1, -1, NULL, pTq->pMetaDB, &pTq->pOffsetStore, 0));
4,893,813✔
486

487
END:
4,893,218✔
488
  if (code != TDB_CODE_SUCCESS) {
4,893,218✔
489
    if (pTq->pExecStore) {
×
490
      tdbTbClose(pTq->pExecStore);
×
491
      pTq->pExecStore = NULL;
×
492
    }
493
    if (pTq->pOffsetStore) {
×
494
      tdbTbClose(pTq->pOffsetStore);
×
495
      pTq->pOffsetStore = NULL;
×
496
    }
497
    tdbClose(pTq->pMetaDB);
×
498
    pTq->pMetaDB = NULL;
×
499
  }
500
  return code;
4,871,388✔
501
}
502

503
static int32_t tqBuildFName(char** data, const char* path, char* name) {
14,678,382✔
504
  int32_t code = 0;
14,678,382✔
505
  int32_t lino = 0;
14,678,382✔
506
  char*   fname = NULL;
14,678,382✔
507
  TSDB_CHECK_NULL(data, code, lino, END, TSDB_CODE_INVALID_MSG);
14,678,382✔
508
  TSDB_CHECK_NULL(path, code, lino, END, TSDB_CODE_INVALID_MSG);
14,678,382✔
509
  TSDB_CHECK_NULL(name, code, lino, END, TSDB_CODE_INVALID_MSG);
14,678,382✔
510
  int32_t len = strlen(path) + strlen(name) + 2;
14,678,382✔
511
  fname = taosMemoryCalloc(1, len);
14,678,382✔
512
  TSDB_CHECK_NULL(fname, code, lino, END, terrno);
14,677,389✔
513
  (void)snprintf(fname, len, "%s%s%s", path, TD_DIRSEP, name);
14,677,389✔
514

515
  *data = fname;
14,677,389✔
516
  fname = NULL;
14,683,875✔
517

518
END:
14,683,875✔
519
  if (code != 0) {
14,683,875✔
520
    tqError("%s failed at %d since %s", __func__, lino, tstrerror(code));
×
521
  }
522
  taosMemoryFree(fname);
14,677,265✔
523
  return code;
14,681,651✔
524
}
525

526
static int32_t tqReplacePath(char** path) {
4,894,182✔
527
  if (path == NULL || *path == NULL) {
4,894,182✔
528
    return TSDB_CODE_INVALID_PARA;
×
529
  }
530
  char*   tpath = NULL;
4,894,665✔
531
  int32_t code = tqBuildFName(&tpath, *path, TQ_SUBSCRIBE_NAME);
4,893,151✔
532
  if (code != 0) {
4,894,047✔
533
    return code;
×
534
  }
535
  taosMemoryFree(*path);
4,894,047✔
536
  *path = tpath;
4,892,862✔
537
  return TDB_CODE_SUCCESS;
4,895,339✔
538
}
539

540
static int32_t tqMetaTransform(STQ* pTq) {
×
541
  if (pTq == NULL) {
×
542
    return TSDB_CODE_INVALID_PARA;
×
543
  }
544
  int32_t code = TDB_CODE_SUCCESS;
×
545
  TDB*    pMetaDB = NULL;
×
546
  TTB*    pExecStore = NULL;
×
547
  char*   offsetNew = NULL;
×
548
  char*   offset = NULL;
×
549
  TQ_ERR_GO_TO_END(tqBuildFName(&offset, pTq->path, TQ_OFFSET_NAME));
×
550

551
  TQ_ERR_GO_TO_END(tdbOpen(pTq->path, 16 * 1024, 1, &pMetaDB, 0, NULL));
×
552
  TQ_ERR_GO_TO_END(tdbTbOpen("tq.db", -1, -1, NULL, pMetaDB, &pExecStore, 0));
×
553

554
  TQ_ERR_GO_TO_END(tqReplacePath(&pTq->path));
×
555
  TQ_ERR_GO_TO_END(tqMetaOpenTdb(pTq));
×
556

557
  TQ_ERR_GO_TO_END(tqMetaTransformInfo(pTq->pMetaDB, pExecStore, pTq->pExecStore));
×
558

559
  TQ_ERR_GO_TO_END(tqBuildFName(&offsetNew, pTq->path, TQ_OFFSET_NAME));
×
560

561
  if (taosCheckExistFile(offset)) {
×
562
    if (taosCopyFile(offset, offsetNew) < 0) {
×
563
      tqError("copy offset file error");
×
564
    } else {
565
      TQ_ERR_GO_TO_END(taosRemoveFile(offset));
×
566
    }
567
  }
568

569
END:
×
570
  taosMemoryFree(offset);
×
571
  taosMemoryFree(offsetNew);
×
572

573
  tdbTbClose(pExecStore);
×
574
  tdbClose(pMetaDB);
×
575
  return code;
×
576
}
577

578
int32_t tqMetaOpen(STQ* pTq) {
4,885,359✔
579
  if (pTq == NULL) {
4,885,359✔
580
    return TSDB_CODE_INVALID_PARA;
×
581
  }
582
  char*   maindb = NULL;
4,885,359✔
583
  char*   offsetNew = NULL;
4,896,078✔
584
  int32_t code = TDB_CODE_SUCCESS;
4,896,604✔
585
  TQ_ERR_GO_TO_END(tqBuildFName(&maindb, pTq->path, TDB_MAINDB_NAME));
4,896,604✔
586
  if (!taosCheckExistFile(maindb)) {
4,892,372✔
587
    TQ_ERR_GO_TO_END(tqReplacePath(&pTq->path));
4,895,785✔
588
    TQ_ERR_GO_TO_END(tqMetaOpenTdb(pTq));
4,895,339✔
589
  } else {
590
    TQ_ERR_GO_TO_END(tqMetaTransform(pTq));
×
591
    TQ_ERR_GO_TO_END(taosRemoveFile(maindb));
×
592
  }
593

594
  TQ_ERR_GO_TO_END(tqBuildFName(&offsetNew, pTq->path, TQ_OFFSET_NAME));
4,894,817✔
595
  if (taosCheckExistFile(offsetNew)) {
4,893,000✔
596
    TQ_ERR_GO_TO_END(tqOffsetRestoreFromFile(pTq, offsetNew));
×
597
    TQ_ERR_GO_TO_END(taosRemoveFile(offsetNew));
×
598
  }
599

600
END:
4,896,459✔
601
  taosMemoryFree(maindb);
4,896,459✔
602
  taosMemoryFree(offsetNew);
4,892,466✔
603
  return code;
4,889,133✔
604
}
605

606
void tqMetaClose(STQ* pTq) {
4,896,920✔
607
  if (pTq == NULL) {
4,896,920✔
608
    return;
×
609
  }
610
  int32_t ret = 0;
4,896,920✔
611
  if (pTq->pExecStore) {
4,896,920✔
612
    tdbTbClose(pTq->pExecStore);
4,896,697✔
613
  }
614
  if (pTq->pOffsetStore) {
4,897,043✔
615
    tdbTbClose(pTq->pOffsetStore);
4,896,920✔
616
  }
617
  tdbClose(pTq->pMetaDB);
4,896,920✔
618
}
STATUS · Troubleshooting · Open an Issue · Sales · Support · CAREERS · ENTERPRISE · START FREE · SCHEDULE DEMO
ANNOUNCEMENTS · TWITTER · TOS & SLA · Supported CI Services · What's a CI service? · Automated Testing

© 2026 Coveralls, Inc