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

taosdata / TDengine / #5119

16 Sep 2026 01:31AM UTC coverage: 73.081% (-0.7%) from 73.739%
#5119

push

travis-ci

jbjia
test(coverage): sync from gitlab

309337 of 423280 relevant lines covered (73.08%)

67769626.34 hits per line

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

82.46
/source/dnode/vnode/src/vnd/vnodeStreamVTable.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
// vnodeStreamVTable.c
17
//
18
// Virtual-table (vtable) reference-chain resolution helpers for the stream
19
// trigger reader path. Extracted from vnodeStream.c to keep that file at a
20
// manageable size. Functions here cover:
21
//   - per-uid VTableInfo packing from resolved chain results
22
//   - throttled vtable-cache recheck hook (WAL meta entry points)
23
//   - single-hop TDMT_VND_VTABLE_REF_RESOLVE server handler
24
//   - multi-hop driver that fans out RPCs and walks the ref chain
25
//   - vchild tag-chain resolver used by executor tag-ref scans
26

27
#include <stdbool.h>
28
#include <stdint.h>
29
#include <taos.h>
30
#include <tdef.h>
31
#include "executor.h"
32
#include "nodes.h"
33
#include "osMemPool.h"
34
#include "osMemory.h"
35
#include "osSemaphore.h"
36
#include "query.h"
37
#include "scalar.h"
38
#include "stream.h"
39
#include "streamReader.h"
40
#include "taosdef.h"
41
#include "taoserror.h"
42
#include "tarray.h"
43
#include "tcommon.h"
44
#include "tdatablock.h"
45
#include "tdb.h"
46
#include "tencode.h"
47
#include "tglobal.h"
48
#include "thash.h"
49
#include "tlist.h"
50
#include "tlockfree.h"
51
#include "tmsg.h"
52
#include "tsimplehash.h"
53
#include "ttypes.h"
54
#include "tutil.h"
55
#include "vnd.h"
56
#include "vnode.h"
57
#include "vnodeInt.h"
58
#include "vnodeStreamVTable.h"
59

60

61
// ---------------------------------------------------------------------------
62
// Block A: per-uid VTableInfo helpers + cache commit
63
// ---------------------------------------------------------------------------
64

65
// Extract tag-typed column ids from the reader's partition-cols node list.
66
// On success returns a fresh SArray<col_id_t> (may be empty); caller frees it.
67
int32_t streamCollectTagCidsFromPartitionCols(SNodeList *partitionCols, SArray **ppTagCids) {
4,269 ✔
68
  *ppTagCids = NULL;
4,269 ✔
69
  SArray *tagCids = taosArrayInit(0, sizeof(col_id_t));
4,269 ✔
70
  if (tagCids == NULL) return terrno;
4,269 ✔
71
  SNode *pNode = NULL;
4,269 ✔
72
  FOREACH(pNode, partitionCols) {
7,547 ✔
73
    if (pNode == NULL || nodeType(pNode) != QUERY_NODE_COLUMN) continue;
3,278 ✔
74
    SColumnNode *c = (SColumnNode *)pNode;
938 ✔
75
    if (c->colType != COLUMN_TYPE_TAG) continue;
938 ✔
76
    col_id_t cid = c->colId;
938 ✔
77
    if (taosArrayPush(tagCids, &cid) == NULL) { taosArrayDestroy(tagCids); return terrno; }
938 ✔
78
  }
79
  *ppTagCids = tagCids;
4,269 ✔
80
  return 0;
4,269 ✔
81
}
82

83
// For a single resolved uid, fill one VTableInfo entry: copy each requested cid's
84
// resolved terminal SColResolveItem into pColRef[i] (hasRef + ref{Db,Table,Col}Name);
85
// id is the virtual cid itself. version is taken from metaReader if available.
86
int32_t streamFillVTableInfoFromResolved(SVnode *pVnode, SStreamTriggerReaderInfo *sStreamReaderInfo,
6,593 ✔
87
                                                int64_t uid, uint64_t gid, int64_t ver, SArray *cids,
88
                                                SVTableResolveResult *pRes, SMetaReader *metaReader,
89
                                                SArray *infos) {
90
  int32_t code = 0;
6,593 ✔
91
  int32_t lino = 0;
6,593 ✔
92
  void   *pTask = sStreamReaderInfo->pTask;
6,593 ✔
93

94
  VTableInfo *vTable = taosArrayReserve(infos, 1);
6,593 ✔
95
  STREAM_CHECK_NULL_GOTO(vTable, terrno);
6,593 ✔
96
  vTable->uid = uid;
6,593 ✔
97
  vTable->gId = gid;
6,593 ✔
98

99
  // Pull schema version + colRef from meta. cids==NULL means "all columns of
100
  // this vtable", in which case we also need me.colRef.pColRef as the iteration
101
  // source. Soft-fail (leave version=0 / nCols=0) if the entry is gone.
102
  int32_t version  = 0;
6,593 ✔
103
  bool    haveMeta = false;
6,593 ✔
104
  code = sStreamReaderInfo->storageApi.metaReaderFn.getTableEntryByVersionUid(metaReader, ver, uid);
6,593 ✔
105
  if (code == 0) {
6,593 ✔
106
    version  = metaReader->me.colRef.version;
6,593 ✔
107
    haveMeta = true;
6,593 ✔
108
  } else {
109
    code = 0;
×
110
  }
111

112
  if (cids == NULL) {
6,593 ✔
113
    // "All columns" mode: enumerate the vtable's own pColRef.
114
    int32_t nAll = haveMeta ? metaReader->me.colRef.nCols : 0;
2,694 ✔
115
    vTable->cols.nCols   = nAll;
2,694 ✔
116
    vTable->cols.version = version;
2,694 ✔
117
    if (nAll > 0) {
2,694 ✔
118
      vTable->cols.pColRef = taosMemoryCalloc(nAll, sizeof(SColRef));
2,694 ✔
119
      STREAM_CHECK_NULL_GOTO(vTable->cols.pColRef, terrno);
2,694 ✔
120
      for (int32_t j = 0; j < nAll; ++j) {
8,961 ✔
121
        col_id_t cid = metaReader->me.colRef.pColRef[j].id;
6,267 ✔
122
        vTable->cols.pColRef[j].id = cid;
6,267 ✔
123
        if (pRes == NULL || pRes->colMap == NULL) continue;
6,267 ✔
124
        SColResolveItem **pp = (SColResolveItem **)tSimpleHashGet(pRes->colMap, &cid, sizeof(cid));
6,267 ✔
125
        if (pp == NULL || *pp == NULL) continue;
6,267 ✔
126
        SColResolveItem *item = *pp;
6,267 ✔
127
        vTable->cols.pColRef[j].hasRef = item->hasRef;
6,267 ✔
128
        if (item->hasRef) {
6,267 ✔
129
          tstrncpy(vTable->cols.pColRef[j].refDbName,    item->refDbName,    TSDB_DB_NAME_LEN);
3,573 ✔
130
          tstrncpy(vTable->cols.pColRef[j].refTableName, item->refTableName, TSDB_TABLE_NAME_LEN);
3,573 ✔
131
          tstrncpy(vTable->cols.pColRef[j].refColName,   item->refColName,   TSDB_COL_NAME_LEN);
3,573 ✔
132
        }
133
      }
134
    }
135
  } else {
136
    int32_t nCids = (int32_t)taosArrayGetSize(cids);
3,899 ✔
137
    vTable->cols.nCols   = nCids;
3,899 ✔
138
    vTable->cols.version = version;
3,899 ✔
139
    vTable->cols.pColRef = taosMemoryCalloc(nCids, sizeof(SColRef));
3,899 ✔
140
    STREAM_CHECK_NULL_GOTO(vTable->cols.pColRef, terrno);
3,899 ✔
141

142
    for (int32_t i = 0; i < nCids; ++i) {
12,032 ✔
143
      col_id_t cid = *(col_id_t *)taosArrayGet(cids, i);
8,133 ✔
144
      vTable->cols.pColRef[i].id = cid;
8,133 ✔
145
      if (pRes == NULL || pRes->colMap == NULL) continue;
8,133 ✔
146
      SColResolveItem **pp = (SColResolveItem **)tSimpleHashGet(pRes->colMap, &cid, sizeof(cid));
8,133 ✔
147
      if (pp == NULL || *pp == NULL) continue;
8,133 ✔
148
      SColResolveItem *item = *pp;
8,133 ✔
149
      vTable->cols.pColRef[i].hasRef = item->hasRef;
8,133 ✔
150
      if (item->hasRef) {
8,133 ✔
151
        tstrncpy(vTable->cols.pColRef[i].refDbName,    item->refDbName,    TSDB_DB_NAME_LEN);
4,033 ✔
152
        tstrncpy(vTable->cols.pColRef[i].refTableName, item->refTableName, TSDB_TABLE_NAME_LEN);
4,033 ✔
153
        tstrncpy(vTable->cols.pColRef[i].refColName,   item->refColName,   TSDB_COL_NAME_LEN);
4,033 ✔
154
      }
155
    }
156
  }
157

158
  if (haveMeta) {
6,593 ✔
159
    tDecoderClear(&metaReader->coder);
6,593 ✔
160
  }
161

162
end:
×
163
  return code;
6,593 ✔
164
}
165

166
int32_t streamCacheCommitResolved(SStreamVTableInfoCache *pCache, bool fullScan,
4,269 ✔
167
                                         SArray *cids, SArray *tagCids, SSHashObj **ppUid2Result) {
168
  int32_t code = 0;
4,269 ✔
169
  if (pCache == NULL || ppUid2Result == NULL || *ppUid2Result == NULL) return TSDB_CODE_INVALID_PARA;
4,269 ✔
170

171
  taosWLockLatch(&pCache->lock);
4,269 ✔
172
  if (fullScan) {
4,269 ✔
173
    TSWAP(pCache->uid2Result, *ppUid2Result);
4,269 ✔
174
  } else {
175
    void *iter = NULL; int32_t it = 0;
×
176
    while ((iter = tSimpleHashIterate(*ppUid2Result, iter, &it)) != NULL) {
×
177
      int64_t                uid = *(int64_t *)tSimpleHashGetKey(iter, NULL);
×
178
      SVTableResolveResult **pSlot = (SVTableResolveResult **)iter;
×
179
      SVTableResolveResult  *r     = *pSlot;
×
180
      if (r == NULL) continue;
×
181
      code = tSimpleHashRemove(pCache->uid2Result, &uid, sizeof(uid));
×
182
      if (code == 0) {
×
183
        code = tSimpleHashPut(pCache->uid2Result, &uid, sizeof(uid), &r, POINTER_BYTES);
×
184
        if (code != 0) {
×
185
          goto _exit;
×
186
        }
187
      }
188
    }
189
  }
190

191
  taosArrayDestroy(pCache->reqColCids);
4,269 ✔
192
  pCache->reqColCids = NULL;
4,269 ✔
193
  if (cids != NULL) {
4,269 ✔
194
    pCache->reqColCids = taosArrayDup(cids, NULL);
2,522 ✔
195
    if (pCache->reqColCids == NULL) { code = terrno; goto _exit; }
2,522 ✔
196
  }
197
  taosArrayDestroy(pCache->reqTagCids);
4,269 ✔
198
  pCache->reqTagCids = NULL;
4,269 ✔
199
  if (tagCids != NULL) {
4,269 ✔
200
    pCache->reqTagCids = taosArrayDup(tagCids, NULL);
4,269 ✔
201
    if (pCache->reqTagCids == NULL) { code = terrno; goto _exit; }
4,269 ✔
202
  }
203
  pCache->lastCheckMs = taosGetTimestampMs();
4,269 ✔
204
  pCache->valid       = true;
4,269 ✔
205

206
_exit:
4,269 ✔
207
  taosWUnLockLatch(&pCache->lock);
4,269 ✔
208
  return code;
4,269 ✔
209
}
210

211
// ---------------------------------------------------------------------------
212
// Block B: throttled vtable cache recheck hook (and its small equality helpers)
213
// ---------------------------------------------------------------------------
214

215
// Compare two SColResolveItem; returns true if they refer to the same terminal column.
216
static bool colResolveItemEqual(const SColResolveItem *a, const SColResolveItem *b) {
106,892 ✔
217
  if (a == NULL && b == NULL) return true;
106,892 ✔
218
  if (a == NULL || b == NULL) return false;
106,892 ✔
219
  if (a->hasRef != b->hasRef) return false;
106,892 ✔
220
  if (!a->hasRef) return true;
106,758 ✔
221
  return strcmp(a->refDbName, b->refDbName) == 0 &&
111,116 ✔
222
         strcmp(a->refTableName, b->refTableName) == 0 &&
111,049 ✔
223
         strcmp(a->refColName, b->refColName) == 0;
55,491 ✔
224
}
225

226
bool tagValueEqual(const STagValue *a, const STagValue *b) {
16,904 ✔
227
  if (a == NULL && b == NULL) return true;
16,904 ✔
228
  if (a == NULL || b == NULL) return false;
16,902 ✔
229
  if (a->type != b->type) return false;
16,898 ✔
230
  if (a->nLen != b->nLen) return false;
16,896 ✔
231
  if (a->nLen == 0) return true;
16,894 ✔
232
  if (a->pData == NULL || b->pData == NULL) return a->pData == b->pData;
16,892 ✔
233
  return memcmp(a->pData, b->pData, a->nLen) == 0;
16,888 ✔
234
}
235

236
// Sliced re-check tuning: every STREAM_VTB_RECHECK_INTERVAL_MS scans at most
237
// STREAM_VTB_RECHECK_SLICE_SIZE uids. A full sweep of N uids therefore takes
238
// roughly ceil(N / SLICE_SIZE) * INTERVAL_MS. With INTERVAL=1000 ms and
239
// SLICE=1000, up to 1000 uids/sec are verified per vnode.
240
#define STREAM_VTB_RECHECK_INTERVAL_MS 1000
241
#define STREAM_VTB_RECHECK_SLICE_SIZE  1000
242
#define STREAM_VTB_RPC_TIMEOUT_MS      30000
243

244
// Throttled hook called at the entry of every WAL meta request.
245
// On tag change: returns TSDB_CODE_STREAM_VTB_TAG_CHANGED so caller bails out fast.
246
// On col-only change: appends affected uids into rsp->tableBlock as TABLE_BLOCK_ADD.
247
// All other cases: returns 0 and lets caller continue normal processing.
248
//
249
// Locking discipline: the resolver round-trip (RPC + tsem2_timewait) is expensive
250
// and MUST run outside the cache W-latch, otherwise WAL meta processing and
251
// the foreground vtable-info request path stall on every recheck tick. The
252
// hook therefore splits work into three phases:
253
//   1) under lock: throttle check + snapshot of reqColCids/reqTagCids and the
254
//      slice uid list, advance the slice cursor, and claim lastCheckMs = now
255
//      so concurrent callers see the throttle and skip;
256
//   2) without lock: streamResolveVTableRefChain over the snapshot;
257
//   3) under lock: diff the resolved result against the live cache and apply
258
//      M1-style per-uid updates / fail-fast on tag changes.
259
// Phase 1 of streamMaybeRecheckVTableCache: take the W-lock, do the double-
260
// checked throttle, refill uidSlice if the cursor wrapped, snapshot the
261
// current slice + reqCol/Tag cids, then advance the cursor and claim the
262
// throttle slot. Drops the lock before returning.
263
//
264
// On success returns 0 and:
265
//   *ppSliceUids != NULL when work should proceed (caller owns and frees).
266
//   *ppSliceUids == NULL when the recheck was throttled or the cache was
267
//                  empty -- caller should treat this as a no-op success.
268
// On failure returns the error code; all out-params are left NULL / 0.
269
static int32_t streamRecheckTakeSlice(SVnode *pVnode, SStreamVTableInfoCache *pCache,
46,166 ✔
270
                                      SStreamTriggerReaderInfo *pInfo,
271
                                      SArray **ppSliceUids, SArray **ppReqColCids,
272
                                      SArray **ppReqTagCids, int32_t *pBegin,
273
                                      int32_t *pEnd, int32_t *pTotal) {
274
  int32_t code       = 0;
46,166 ✔
275
  SArray *sliceUids  = NULL;
46,166 ✔
276
  SArray *reqColCids = NULL;
46,166 ✔
277
  SArray *reqTagCids = NULL;
46,166 ✔
278
  int32_t begin = 0, end = 0, total = 0;
46,166 ✔
279

280
  *ppSliceUids  = NULL;
46,166 ✔
281
  *ppReqColCids = NULL;
46,166 ✔
282
  *ppReqTagCids = NULL;
46,166 ✔
283
  *pBegin = *pEnd = *pTotal = 0;
46,166 ✔
284

285
  taosWLockLatch(&pCache->lock);
46,166 ✔
286

287
  // Double-checked throttle: concurrent caller may have just done a sweep.
288
  int64_t now = taosGetTimestampMs();
46,230 ✔
289
  if (now - pCache->lastCheckMs < STREAM_VTB_RECHECK_INTERVAL_MS) {
46,230 ✔
290
    taosWUnLockLatch(&pCache->lock);
4,721 ✔
291
    return 0;
4,721 ✔
292
  }
293

294
  // Refill uidSlice whenever the cursor wraps to 0 so newly registered vtables
295
  // (not yet in uid2Result) are picked up by the next sweep.
296
  if (pCache->sliceCursor == 0) {
41,509 ✔
297
    taosArrayClear(pCache->uidSlice);
41,445 ✔
298
    SArray *pTableListArray = qStreamGetTableArrayList(pInfo);
41,381 ✔
299
    if (pTableListArray == NULL) {
41,509 ✔
300
      taosWUnLockLatch(&pCache->lock);
×
301
      return terrno;
×
302
    }
303
    int32_t nAll = (int32_t)taosArrayGetSize(pTableListArray);
41,509 ✔
304
    for (int32_t i = 0; i < nAll; ++i) {
102,908 ✔
305
      SStreamTableKeyInfo *pKey = taosArrayGetP(pTableListArray, i);
61,399 ✔
306
      if (pKey == NULL || pKey->markedDeleted) continue;
61,399 ✔
307
      if (taosArrayPush(pCache->uidSlice, &pKey->uid) == NULL) {
122,798 ✔
308
        code = terrno;
×
309
        taosArrayDestroyP(pTableListArray, taosMemFree);
×
310
        goto _unlock;
×
311
      }
312
    }
313
    taosArrayDestroyP(pTableListArray, taosMemFree);
41,509 ✔
314
  }
315

316
  total = (int32_t)taosArrayGetSize(pCache->uidSlice);
41,573 ✔
317
  if (total == 0) {
41,509 ✔
318
    pCache->lastCheckMs = taosGetTimestampMs();
8,791 ✔
319
    taosWUnLockLatch(&pCache->lock);
8,791 ✔
320
    stDebug("vgId:%d %s skip: cache empty", TD_VID(pVnode), __func__);
8,791 ✔
321
    return 0;
8,791 ✔
322
  }
323

324
  begin = pCache->sliceCursor;
32,718 ✔
325
  end   = TMIN(begin + STREAM_VTB_RECHECK_SLICE_SIZE, total);
32,718 ✔
326
  sliceUids = taosArrayInit(end - begin, sizeof(int64_t));
32,718 ✔
327
  if (sliceUids == NULL) { code = terrno; goto _unlock; }
32,718 ✔
328
  for (int32_t i = begin; i < end; ++i) {
94,117 ✔
329
    if (taosArrayPush(sliceUids, taosArrayGet(pCache->uidSlice, i)) == NULL) {
122,798 ✔
330
      code = terrno;
×
331
      goto _unlock;
×
332
    }
333
  }
334

335
  if (pCache->reqColCids != NULL) {
32,718 ✔
336
    reqColCids = taosArrayDup(pCache->reqColCids, NULL);
20,949 ✔
337
    if (reqColCids == NULL) { code = terrno; goto _unlock; }
20,949 ✔
338
  }
339
  if (pCache->reqTagCids != NULL) {
32,718 ✔
340
    reqTagCids = taosArrayDup(pCache->reqTagCids, NULL);
32,718 ✔
341
    if (reqTagCids == NULL) { code = terrno; goto _unlock; }
32,718 ✔
342
  }
343

344
  // Advance cursor and claim the throttle slot so concurrent callers skip.
345
  pCache->sliceCursor = (end >= total) ? 0 : end;
32,718 ✔
346
  pCache->lastCheckMs = taosGetTimestampMs();
32,718 ✔
347

348
_unlock:
32,718 ✔
349
  taosWUnLockLatch(&pCache->lock);
32,718 ✔
350
  if (code != 0) {
32,718 ✔
351
    taosArrayDestroy(sliceUids);
×
352
    taosArrayDestroy(reqColCids);
×
353
    taosArrayDestroy(reqTagCids);
×
354
    return code;
×
355
  }
356
  *ppSliceUids  = sliceUids;
32,718 ✔
357
  *ppReqColCids = reqColCids;
32,718 ✔
358
  *ppReqTagCids = reqTagCids;
32,718 ✔
359
  *pBegin = begin; *pEnd = end; *pTotal = total;
32,718 ✔
360
  return 0;
32,718 ✔
361
}
362

363
// Phase 3 per-uid body of streamMaybeRecheckVTableCache: diff the freshly
364
// resolved newRes against the cached oldRes for one uid; on tag change return
365
// TSDB_CODE_STREAM_VTB_TAG_CHANGED so the caller bails out of the loop; on
366
// col change append uid to changedUids; otherwise replace the cache entry
367
// with newRes (ownership transferred from uid2Result).
368
//
369
// Must be called with pCache->lock held in write mode (caller's responsibility).
370
static int32_t streamRecheckDiffAndApplyOneUid(SVnode *pVnode, SStreamVTableInfoCache *pCache,
61,265 ✔
371
                                               int64_t uid, SSHashObj *uid2Result,
372
                                               SArray *changedUids) {
373
  SVTableResolveResult **ppNew = (SVTableResolveResult **)tSimpleHashGet(uid2Result, &uid, sizeof(uid));
61,265 ✔
374
  SVTableResolveResult **ppOld = (SVTableResolveResult **)tSimpleHashGet(pCache->uid2Result, &uid, sizeof(uid));
61,265 ✔
375
  SVTableResolveResult  *newRes = (ppNew == NULL) ? NULL : *ppNew;
61,265 ✔
376
  SVTableResolveResult  *oldRes = (ppOld == NULL) ? NULL : *ppOld;
61,265 ✔
377

378
  // uid skipped by resolver (top-level vtable dropped, H2 fallback) -> drop from cache.
379
  if (newRes == NULL) {
61,265 ✔
380
    if (oldRes != NULL) {
9,546 ✔
381
      stDebug("vgId:%d %s uid dropped: uid=%" PRId64, TD_VID(pVnode), __func__, uid);
4,662 ✔
382
      int32_t rc = tSimpleHashRemove(pCache->uid2Result, &uid, sizeof(uid));
4,662 ✔
383
      if (rc != 0) {
4,662 ✔
384
        stWarn("vgId:%d %s remove uid=%" PRId64 " from cache failed: 0x%x",
×
385
               TD_VID(pVnode), __func__, uid, rc);
386
      }
387
    }
388
    return 0;
9,546 ✔
389
  }
390

391
  // Tag diff -- any tag change is fatal.
392
  bool tagChanged = false;
51,719 ✔
393
  if (oldRes != NULL && oldRes->tagMap != NULL) {
51,719 ✔
394
    void *it2 = NULL; int32_t i2 = 0;
47,649 ✔
395
    while ((it2 = tSimpleHashIterate(oldRes->tagMap, it2, &i2)) != NULL) {
64,466 ✔
396
      col_id_t   cid  = *(col_id_t *)tSimpleHashGetKey(it2, NULL);
16,884 ✔
397
      STagValue *oldV = *(STagValue **)it2;
16,884 ✔
398
      STagValue **ppNewV = (newRes->tagMap == NULL) ? NULL :
16,884 ✔
399
                           (STagValue **)tSimpleHashGet(newRes->tagMap, &cid, sizeof(cid));
16,884 ✔
400
      STagValue  *newV   = (ppNewV == NULL) ? NULL : *ppNewV;
16,884 ✔
401
      if (!tagValueEqual(oldV, newV)) {
16,884 ✔
402
        stDebug("vgId:%d %s tag changed: uid=%" PRId64 " cid=%d", TD_VID(pVnode), __func__,
67 ✔
403
                uid, (int32_t)cid);
404
        tagChanged = true;
67 ✔
405
        break;
67 ✔
406
      }
407
    }
408
  }
409
  if (tagChanged) return TSDB_CODE_STREAM_VTB_TAG_CHANGED;
51,719 ✔
410

411
  // Col diff -- collect uids that need re-publication.
412
  bool colChanged = false;
51,652 ✔
413
  if (oldRes != NULL && oldRes->colMap != NULL) {
51,652 ✔
414
    void *it2 = NULL; int32_t i2 = 0;
47,582 ✔
415
    while ((it2 = tSimpleHashIterate(oldRes->colMap, it2, &i2)) != NULL) {
154,273 ✔
416
      col_id_t          cid  = *(col_id_t *)tSimpleHashGetKey(it2, NULL);
106,892 ✔
417
      SColResolveItem  *oldI = *(SColResolveItem **)it2;
106,892 ✔
418
      SColResolveItem **ppNewI = (newRes->colMap == NULL) ? NULL :
106,892 ✔
419
                                 (SColResolveItem **)tSimpleHashGet(newRes->colMap, &cid, sizeof(cid));
106,892 ✔
420
      SColResolveItem  *newI   = (ppNewI == NULL) ? NULL : *ppNewI;
106,892 ✔
421
      if (!colResolveItemEqual(oldI, newI)) { colChanged = true; break; }
106,892 ✔
422
    }
423
  }
424
  if (colChanged) {
51,652 ✔
425
    if (taosArrayPush(changedUids, &uid) == NULL) return terrno;
201 ✔
426
  }
427

428
  // Replace cache entry with the freshly resolved result; transfer ownership.
429
  if (tSimpleHashPut(pCache->uid2Result, &uid, sizeof(uid), &newRes, POINTER_BYTES) != 0) {
51,652 ✔
430
    return terrno;
×
431
  }
432

433
  streamVTableResolveResultDestroy(&oldRes);
51,652 ✔
434
  *ppNew = NULL;
51,652 ✔
435
  return 0;
51,652 ✔
436
}
437

438
int32_t streamMaybeRecheckVTableCache(SVnode *pVnode, SStreamTriggerReaderInfo *pInfo,
1,931,077 ✔
439
                                             int64_t walVer, SSTriggerWalNewRsp *pRsp) {
440
  if (pInfo == NULL || pInfo->vtbCache == NULL || !pInfo->vtbCache->valid) {
1,931,077 ✔
441
    return 0;
1,882,243 ✔
442
  }
443
  SStreamVTableInfoCache *pCache = pInfo->vtbCache;
46,230 ✔
444
  // Throttle check is done under lock in streamRecheckTakeSlice (Phase 1).
445

446
  int32_t    code         = 0;
46,230 ✔
447
  SArray    *sliceUids    = NULL;
46,230 ✔
448
  SArray    *reqColCids   = NULL;
46,166 ✔
449
  SArray    *reqTagCids   = NULL;
46,166 ✔
450
  SArray    *changedUids  = NULL;
46,166 ✔
451
  SSHashObj *uid2Result   = NULL;
46,166 ✔
452
  int32_t    begin = 0, end = 0, total = 0;
46,166 ✔
453

454
  // ---- Phase 1: snapshot under lock ----
455
  code = streamRecheckTakeSlice(pVnode, pCache, pInfo, &sliceUids, &reqColCids, &reqTagCids,
46,166 ✔
456
                                &begin, &end, &total);
457
  if (code != 0) goto _cleanup;
46,163 ✔
458
  if (sliceUids == NULL) return 0;  // throttled or empty cache
46,163 ✔
459

460
  stDebug("vgId:%d %s walVer=%" PRId64 " total=%d slice=[%d,%d)",
32,651 ✔
461
          TD_VID(pVnode), __func__, walVer, total, begin, end);
462

463
  // ---- Phase 2: resolver round-trip, no lock held ----
464
  code = streamResolveVTableRefChain(pVnode, pCache, pInfo, walVer, sliceUids,
32,651 ✔
465
                                     reqColCids, reqTagCids, &uid2Result);
466
  if (code != 0) goto _cleanup;
32,718 ✔
467

468
  changedUids = taosArrayInit(0, sizeof(int64_t));
32,584 ✔
469
  if (changedUids == NULL) { code = terrno; goto _cleanup; }
32,584 ✔
470

471
  // ---- Phase 3: diff + apply under lock ----
472
  taosWLockLatch(&pCache->lock);
32,584 ✔
473
  for (int32_t i = 0; i < (int32_t)taosArrayGetSize(sliceUids); ++i) {
93,782 ✔
474
    int64_t uid = *(int64_t *)taosArrayGet(sliceUids, i);
61,265 ✔
475
    code = streamRecheckDiffAndApplyOneUid(pVnode, pCache, uid, uid2Result, changedUids);
61,265 ✔
476
    if (code != 0) break;
61,265 ✔
477
  }
478
  pCache->lastCheckMs = taosGetTimestampMs();
32,584 ✔
479
  taosWUnLockLatch(&pCache->lock);
32,584 ✔
480

481
_cleanup:
32,718 ✔
482
  tSimpleHashCleanup(uid2Result);
32,718 ✔
483
  taosArrayDestroy(sliceUids);
32,718 ✔
484
  taosArrayDestroy(reqColCids);
32,718 ✔
485
  taosArrayDestroy(reqTagCids);
32,718 ✔
486

487
  if (code == TSDB_CODE_STREAM_VTB_TAG_CHANGED) {
32,718 ✔
488
    stWarn("vgId:%d %s tag changed, abort fast walVer=%" PRId64, TD_VID(pVnode), __func__, walVer);
67 ✔
489
    taosArrayDestroy(changedUids);
67 ✔
490
    return code;
67 ✔
491
  }
492
  if (code != 0) {
32,651 ✔
493
    stError("vgId:%d %s recheck failed since %s", TD_VID(pVnode), __func__, tstrerror(code));
134 ✔
494
    taosArrayDestroy(changedUids);
134 ✔
495
    return code;
134 ✔
496
  }
497
  if (pRsp != NULL && changedUids != NULL && taosArrayGetSize(changedUids) > 0) {
32,517 ✔
498
    int32_t rc = addUidListToBlock(changedUids, &pRsp->tableBlock, walVer, &pRsp->totalRows, TABLE_BLOCK_ADD);
201 ✔
499
    stDebug("vgId:%d %s appended %d changed uids walVer=%" PRId64, TD_VID(pVnode), __func__,
201 ✔
500
            (int32_t)taosArrayGetSize(changedUids), walVer);
501
    if (rc != 0) { taosArrayDestroy(changedUids); return rc; }
201 ✔
502
  }
503
  taosArrayDestroy(changedUids);
32,517 ✔
504
  return 0;
32,517 ✔
505
}
506

507
// ---------------------------------------------------------------------------
508
// Block C: vtable chain resolution
509
//
510
// This block has THREE sub-sections that should be kept distinct. Do NOT call
511
// internal helpers across sub-section boundaries; the only legal entry points
512
// across sub-sections are the public top-level functions listed below.
513
//
514
//   C1 (server)   - TDMT_VND_VTABLE_REF_RESOLVE single-hop server handler.
515
//                   Public entry: vnodeProcessVTableRefResolveReq.
516
//                   Statics:      vnodeFindVTableColRef, vnodeResolveTableGroup,
517
//                                 vnodeResolveOneHop (also reused by C2 local fast-path),
518
//                                 vnodeFillTagValueFromChild.
519
//
520
//   C2 (driver)   - multi-hop client driver: fans out per-vgId RPCs, walks the
521
//                   chain, integrates with reader tblRefCache / dbVgInfo cache.
522
//                   Public entry: streamResolveVTableRefChain.
523
//
524
//   C3 (executor) - vchild-tag chain helper used by executor tag-ref scans.
525
//                   Public entry: vnodeResolveVTableTagChain.
526
//
527
// If a new RPC type is added in the future, prefer splitting C1 / C2 into
528
// dedicated files (vnodeStreamVTableRpc.c / vnodeStreamVTableDriver.c) over
529
// growing this file further.
530
// ---------------------------------------------------------------------------
531

532
// ---------------------------------------------------------------------------
533
// C1: single-hop RPC server (TDMT_VND_VTABLE_REF_RESOLVE)
534
// ---------------------------------------------------------------------------
535

536
// ============================================================================
537
// TDMT_VND_VTABLE_REF_RESOLVE — single-hop chain resolver for vtable references.
538
//
539
// Caller (driver A in stream-trigger reader info path) groups per-batch refs by
540
// vgId and sends one request per vgId. Each item carries (kind, refDb, refTbl,
541
// refCol). For every item we do exactly ONE hop:
542
//   - table not on this vnode               -> r.code = STREAM_VTB_REF_TABLE_NOT_EXIST
543
//   - column/tag name not found             -> r.code = STREAM_VTB_REF_COL_NOT_EXIST
544
//   - vtable + COL + hasRef                 -> r.terminated = false, r.nextRef = stored ref
545
//   - vtable + COL + !hasRef                -> r.terminated = true,  r.nextRef.hasRef = false
546
//                                              (terminal triple is meaningless; signals NULL value)
547
//   - vtable + TAG + hasRef                 -> r.terminated = false, r.nextRef = stored ref
548
//   - vchild + TAG + !hasRef                -> r.terminated = true,  r.nextRef.hasRef = false,
549
//                                              r.tagType/tagLen/tagData filled from local STag
550
//   - vnormal+ TAG + !hasRef                -> r.code = STREAM_VTB_REF_COL_NOT_EXIST
551
//                                              (normal vtable has no tag concept)
552
//   - physical table  + COL kind            -> r.terminated = true,  r.nextRef = current triple
553
//   - child table     + TAG kind            -> r.terminated = true,  r.tagType/tagLen/tagData filled
554
//   - normal table    + TAG kind            -> r.code = STREAM_VTB_REF_COL_NOT_EXIST
555
// Per-item errors never abort the batch — they are reported in r.code.
556
// ============================================================================
557

558
// Reads a tag's constant value from a (virtual or physical) child table entry.
559
// The stable schema (for type/colId lookup) is fetched here under META_READER_LOCK;
560
// the child entry is provided by the caller.
561
//   pVnode      : owning vnode
562
//   pChildEntry : decoded entry whose type is *_CHILD_TABLE (carries ctbEntry.suid/pTags)
563
//   tagColName  : tag name on the stable (vchild's SColRef.colName matches stable tag name
564
//                 by build-time convention)
565
// Outputs:
566
//   *outType    : tag SDataType
567
//   *outLen     : payload length; 0 when tag absent on this child
568
//   *outData    : newly allocated buffer (caller frees); NULL when *outLen==0
569
// Returns:
570
//   0                                          success (incl. "tag absent")
571
//   TSDB_CODE_STREAM_VTB_REF_TABLE_NOT_EXIST   suid not present on this vnode
572
//   TSDB_CODE_STREAM_VTB_REF_COL_NOT_EXIST     tag name not in stable schema
573
//   terrno                                     OOM
574
// Internal helper: read a constant tag value from a virtual child table.
575
// Tag is located in the parent stable's schemaTag by either colId (preferred when > 0)
576
// or by colName (fallback). vtable on-disk SColRef does not persist colName, so callers
577
// holding only a SColRef entry must pass the cid.
578
static int32_t streamReadChildTagConstValueImpl(SVnode *pVnode, const SMetaEntry *pChildEntry,
24,236 ✔
579
                                                col_id_t tagColId, const char *tagColName,
580
                                                int8_t *outType, int32_t *outLen, char **outData) {
581
  SMetaReader stb  = {0};
24,236 ✔
582
  int32_t     code = 0;
24,236 ✔
583
  *outType = 0;
24,236 ✔
584
  *outLen  = 0;
24,236 ✔
585
  *outData = NULL;
24,236 ✔
586

587
  metaReaderDoInit(&stb, pVnode->pMeta, META_READER_LOCK, 0);
24,236 ✔
588
  if (metaReaderGetTableEntryByUid(&stb, pChildEntry->ctbEntry.suid) != 0) {
24,236 ✔
589
    code = TSDB_CODE_STREAM_VTB_REF_TABLE_NOT_EXIST;
×
590
    goto _end;
×
591
  }
592

593
  SSchemaWrapper *pSW = &stb.me.stbEntry.schemaTag;
24,177 ✔
594
  SSchema        *pTagSchema = NULL;
24,177 ✔
595
  for (int32_t i = 0; i < pSW->nCols; ++i) {
25,785 ✔
596
    if (tagColId > 0) {
25,785 ✔
597
      if (pSW->pSchema[i].colId == tagColId) {
15,534 ✔
598
        pTagSchema = &pSW->pSchema[i];
15,534 ✔
599
        break;
15,534 ✔
600
      }
601
    } else if (tagColName != NULL &&
10,251 ✔
602
               strncmp(pSW->pSchema[i].name, tagColName, TSDB_COL_NAME_LEN) == 0) {
10,251 ✔
603
      pTagSchema = &pSW->pSchema[i];
8,643 ✔
604
      break;
8,643 ✔
605
    }
606
  }
607
  if (pTagSchema == NULL) {
24,177 ✔
608
    code = TSDB_CODE_STREAM_VTB_REF_COL_NOT_EXIST;
×
609
    goto _end;
×
610
  }
611

612
  *outType = pTagSchema->type;
24,177 ✔
613

614
  STag   *pTag  = (STag *)pChildEntry->ctbEntry.pTags;
24,177 ✔
615
  STagVal tv    = {.cid = pTagSchema->colId, .type = pTagSchema->type};
24,177 ✔
616
  bool    found = (pTag != NULL) && tTagGet(pTag, &tv);
24,177 ✔
617
  if (!found) {
24,236 ✔
618
    // tag has no value on this child: outLen=0 / outData=NULL
619
    goto _end;
×
620
  }
621

622
  if (IS_VAR_DATA_TYPE(pTagSchema->type)) {
24,236 ✔
623
    *outLen = (int32_t)tv.nData;
283 ✔
624
    if (*outLen > 0) {
283 ✔
625
      *outData = taosMemoryMalloc(*outLen);
165 ✔
626
      if (*outData == NULL) { code = terrno; goto _end; }
165 ✔
627
      memcpy(*outData, tv.pData, *outLen);
165 ✔
628
    }
629
  } else {
630
    *outLen  = (int32_t)tDataTypes[pTagSchema->type].bytes;
23,953 ✔
631
    *outData = taosMemoryMalloc(*outLen);
23,953 ✔
632
    if (*outData == NULL) { code = terrno; goto _end; }
24,071 ✔
633
    memcpy(*outData, &tv.i64, *outLen);
24,071 ✔
634
  }
635

636
_end:
24,354 ✔
637
  metaReaderClear(&stb);
24,118 ✔
638
  return code;
24,177 ✔
639
}
640

641
// Look up by name (used by request-driven path where wire holds refColName).
642
static int32_t streamReadChildTagConstValue(SVnode *pVnode, const SMetaEntry *pChildEntry,
8,643 ✔
643
                                            const char *tagColName, int8_t *outType,
644
                                            int32_t *outLen, char **outData) {
645
  return streamReadChildTagConstValueImpl(pVnode, pChildEntry, 0, tagColName,
8,643 ✔
646
                                          outType, outLen, outData);
647
}
648

649
// Look up by cid (used by local seed where SColRef.colName is not persisted).
650
static int32_t streamReadChildTagConstValueByCid(SVnode *pVnode, const SMetaEntry *pChildEntry,
15,475 ✔
651
                                                 col_id_t tagColId, int8_t *outType,
652
                                                 int32_t *outLen, char **outData) {
653
  return streamReadChildTagConstValueImpl(pVnode, pChildEntry, tagColId, NULL,
15,475 ✔
654
                                          outType, outLen, outData);
655
}
656

657
static int32_t vnodeFillTagValueFromChild(SVnode *pVnode, const SMetaEntry *pChildEntry,
8,643 ✔
658
                                          const char *tagColName, SVTableRefResolveRspItem *r) {
659
  r->terminated = true;
8,643 ✔
660
  int32_t code = streamReadChildTagConstValue(pVnode, pChildEntry, tagColName,
8,643 ✔
661
                                              &r->tagType, &r->tagLen, &r->tagData);
662
  vDebug("vgId:%d %s tag=%s code=0x%x type=%d len=%d", TD_VID(pVnode), __func__, tagColName, code,
8,643 ✔
663
         r->tagType, r->tagLen);
664
  if (code == TSDB_CODE_STREAM_VTB_REF_TABLE_NOT_EXIST ||
8,643 ✔
665
      code == TSDB_CODE_STREAM_VTB_REF_COL_NOT_EXIST) {
666
    // per-item soft error: surface via r->code, do not abort the batch.
667
    r->code = code;
×
668
    return 0;
×
669
  }
670
  return code;
8,643 ✔
671
}
672

673
// Lookup the SColRef on a vtable for (kind, colName), returning a pointer into
674
// the vtable's pColRef/pTagRef array (or NULL when the column name is not in
675
// the schema — this is a per-item soft miss, not a function error).
676
//
677
// For VIRTUAL_CHILD_TABLE the parent stable's schema is needed to translate
678
// colName -> cid. If the caller has already opened the parent stable entry
679
// (e.g. to share it across multiple columns of the same vchild), it can pass
680
// pStbEntry to avoid the extra meta read; otherwise pass NULL and the helper
681
// will open & close a temporary reader internally.
682
//
683
// Returns: 0 on success (*ppFoundRef set, may be NULL).
684
//          TSDB_CODE_STREAM_VTB_REF_TABLE_NOT_EXIST when pStbEntry==NULL and
685
//          the parent stable cannot be opened — the caller should surface this
686
//          via the per-item rsp.code.
687
static int32_t vnodeFindVTableColRef(SVnode *pVnode, const SMetaEntry *pVtbEntry,
23,718 ✔
688
                                     const SMetaEntry *pStbEntry, EStreamVRefKind kind,
689
                                     const char *colName, SColRef **ppFoundRef) {
690
  *ppFoundRef = NULL;
23,718 ✔
691
  if (pVtbEntry->type != TSDB_VIRTUAL_NORMAL_TABLE &&
23,718 ✔
692
      pVtbEntry->type != TSDB_VIRTUAL_CHILD_TABLE) {
18,425 ✔
693
    return 0;
×
694
  }
695

696
  const SColRefWrapper *pWrap = &pVtbEntry->colRef;
23,718 ✔
697
  SColRef              *pArr  = (kind == STREAM_VREF_KIND_TAG) ? pWrap->pTagRef : pWrap->pColRef;
23,718 ✔
698
  int32_t               nArr  = (kind == STREAM_VREF_KIND_TAG) ? pWrap->nTagRefs : pWrap->nCols;
23,718 ✔
699

700
  SMetaReader           tmpStb    = {0};
23,718 ✔
701
  bool                  tmpInited = false;
23,718 ✔
702
  const SSchemaWrapper *pSW       = NULL;
23,718 ✔
703

704
  if (pVtbEntry->type == TSDB_VIRTUAL_NORMAL_TABLE) {
23,718 ✔
705
    // Normal vtable: schema lives on the vtable entry itself.
706
    pSW = &pVtbEntry->ntbEntry.schemaRow;
5,293 ✔
707
  } else if (pStbEntry != NULL) {
18,425 ✔
708
    // Caller pre-loaded parent stable — reuse it.
709
    pSW = (kind == STREAM_VREF_KIND_TAG) ? &pStbEntry->stbEntry.schemaTag
6,365 ✔
710
                                         : &pStbEntry->stbEntry.schemaRow;
6,365 ✔
711
  } else {
712
    // Open parent stable on the fly.
713
    metaReaderDoInit(&tmpStb, pVnode->pMeta, META_READER_LOCK, 0);
12,060 ✔
714
    if (metaReaderGetTableEntryByUid(&tmpStb, pVtbEntry->ctbEntry.suid) != 0) {
12,060 ✔
715
      vDebug("vgId:%d %s parent stable not found: suid=%" PRId64, TD_VID(pVnode), __func__,
×
716
             pVtbEntry->ctbEntry.suid);
717
      metaReaderClear(&tmpStb);
×
718
      return TSDB_CODE_STREAM_VTB_REF_TABLE_NOT_EXIST;
×
719
    }
720
    tmpInited = true;
12,060 ✔
721
    pSW = (kind == STREAM_VREF_KIND_TAG) ? &tmpStb.me.stbEntry.schemaTag
12,060 ✔
722
                                         : &tmpStb.me.stbEntry.schemaRow;
723
  }
724

725
  // Step 1: resolve colName -> cid using the chosen schema wrapper.
726
  col_id_t targetCid = 0;
23,718 ✔
727
  bool     cidFound  = false;
23,718 ✔
728
  for (int32_t k = 0; pSW != NULL && k < pSW->nCols; ++k) {
38,659 ✔
729
    if (strncmp(pSW->pSchema[k].name, colName, TSDB_COL_NAME_LEN) == 0) {
38,592 ✔
730
      targetCid = pSW->pSchema[k].colId;
23,651 ✔
731
      cidFound  = true;
23,651 ✔
732
      break;
23,651 ✔
733
    }
734
  }
735

736
  // Step 2: scan the vtable's ref array for that cid.
737
  if (cidFound) {
23,718 ✔
738
    for (int32_t j = 0; j < nArr && pArr != NULL; ++j) {
38,458 ✔
739
      if (pArr[j].id == targetCid) {
38,458 ✔
740
        *ppFoundRef = &pArr[j];
23,651 ✔
741
        break;
23,651 ✔
742
      }
743
    }
744
  }
745

746
  if (tmpInited) metaReaderClear(&tmpStb);
23,718 ✔
747
  return 0;
23,718 ✔
748
}
749

750
// Fill a SVTableRefResolveRspItem based on the resolved SColRef (or lack thereof).
751
// Handles all vtable/physical-table × COL/TAG × hasRef/!hasRef combinations in
752
// one place so that both vnodeResolveTableGroup and vnodeResolveOneHop share the
753
// same branching logic without duplication.
754
//
755
// Parameters:
756
//   pVnode    - vnode handle (for vnodeFillTagValueFromChild)
757
//   pEntry   - the meta entry of the table being resolved
758
//   pFound   - resolved SColRef, may be NULL (col not found)
759
//   kind     - STREAM_VREF_KIND_COL or STREAM_VREF_KIND_TAG
760
//   dbName   - database name for the physical-table terminal triple
761
//   tableName- table name for the physical-table terminal triple
762
//   colName  - column name (for tag value lookup)
763
//   isVtable - whether the table is a virtual table
764
//   r        - output response item (caller zeroes it before call)
765
//
766
// Returns 0 on success, or an error code for fatal failures (e.g. OOM in tag
767
// fill). Logical resolution errors (col-not-exist etc.) are written into r->code
768
// and the function still returns 0.
769
static int32_t vnodeFillResolveRspFromColRef(SVnode *pVnode, const SMetaEntry *pEntry,
98,651 ✔
770
                                             SColRef *pFound, int8_t kind,
771
                                             const char *dbName, const char *tableName,
772
                                             const char *colName, bool isVtable,
773
                                             SVTableRefResolveRspItem *r) {
774
  if (isVtable) {
98,651 ✔
775
    if (pFound == NULL) {
23,718 ✔
776
      r->code = TSDB_CODE_STREAM_VTB_REF_COL_NOT_EXIST;
67 ✔
777
      return 0;
67 ✔
778
    }
779
    if (pFound->hasRef) {
23,651 ✔
780
      r->terminated     = false;
23,651 ✔
781
      r->nextRef.kind   = kind;
23,651 ✔
782
      r->nextRef.hasRef = true;
23,651 ✔
783
      tstrncpy(r->nextRef.refDbName,    pFound->refDbName,    TSDB_DB_NAME_LEN);
23,651 ✔
784
      tstrncpy(r->nextRef.refTableName, pFound->refTableName, TSDB_TABLE_NAME_LEN);
23,651 ✔
785
      tstrncpy(r->nextRef.refColName,   pFound->refColName,   TSDB_COL_NAME_LEN);
23,651 ✔
786
      return 0;
23,651 ✔
787
    }
788
    // !hasRef
789
    if (kind == STREAM_VREF_KIND_TAG) {
×
790
      if (pEntry->type != TSDB_VIRTUAL_CHILD_TABLE) {
×
791
        r->code = TSDB_CODE_STREAM_VTB_REF_COL_NOT_EXIST;
×
792
        return 0;
×
793
      }
794
      r->nextRef.kind            = STREAM_VREF_KIND_TAG;
×
795
      r->nextRef.hasRef          = false;
×
796
      r->nextRef.refDbName[0]    = '\0';
×
797
      r->nextRef.refTableName[0] = '\0';
×
798
      r->nextRef.refColName[0]   = '\0';
×
799
      int32_t rc = vnodeFillTagValueFromChild(pVnode, pEntry, colName, r);
×
800
      if (rc != 0) r->code = rc;
×
801
      return 0;
×
802
    }
803
    // STREAM_VREF_KIND_COL on vtable with NULL ref: terminal empty
804
    r->terminated              = true;
×
805
    r->nextRef.kind            = STREAM_VREF_KIND_COL;
×
806
    r->nextRef.hasRef          = false;
×
807
    r->nextRef.refDbName[0]    = '\0';
×
808
    r->nextRef.refTableName[0] = '\0';
×
809
    r->nextRef.refColName[0]   = '\0';
×
810
    return 0;
×
811
  }
812

813
  // Physical table
814
  if (kind == STREAM_VREF_KIND_COL) {
74,933 ✔
815
    r->terminated     = true;
66,290 ✔
816
    r->nextRef.kind   = STREAM_VREF_KIND_COL;
66,290 ✔
817
    r->nextRef.hasRef = true;
66,290 ✔
818
    tstrncpy(r->nextRef.refDbName,    dbName,    TSDB_DB_NAME_LEN);
66,290 ✔
819
    tstrncpy(r->nextRef.refTableName, tableName, TSDB_TABLE_NAME_LEN);
66,290 ✔
820
    tstrncpy(r->nextRef.refColName,   colName,   TSDB_COL_NAME_LEN);
66,290 ✔
821
    return 0;
66,290 ✔
822
  }
823

824
  // TAG on physical table: only child table carries tag values
825
  if (pEntry->type != TSDB_CHILD_TABLE) {
8,643 ✔
826
    r->code = TSDB_CODE_STREAM_VTB_REF_COL_NOT_EXIST;
×
827
    return 0;
×
828
  }
829
  r->nextRef.kind   = STREAM_VREF_KIND_TAG;
8,643 ✔
830
  r->nextRef.hasRef = false;
8,643 ✔
831
  int32_t rc = vnodeFillTagValueFromChild(pVnode, pEntry, colName, r);
8,643 ✔
832
  if (rc != 0) r->code = rc;
8,643 ✔
833
  return 0;
8,643 ✔
834
}
835

836
// Batch-resolve multiple columns within the same table. Opens meta once for the
837
// table, then resolves each (colName, kind) pair against the same metadata.
838
// Results are appended to pRspItems in the same order as pCols.
839
static int32_t vnodeResolveTableGroup(SVnode *pVnode, const char *dbName, const char *tableName,
39,044 ✔
840
                                      SArray *pCols, SArray *pRspItems) {
841
  SMetaReader mr   = {0};
39,044 ✔
842
  int32_t     code = 0;
39,044 ✔
843
  int32_t     nCols = (pCols != NULL) ? (int32_t)taosArrayGetSize(pCols) : 0;
39,044 ✔
844

845
  vDebug("vgId:%d %s enter: db=%s table=%s nCols=%d", TD_VID(pVnode), __func__, dbName, tableName, nCols);
39,044 ✔
846

847
  metaReaderDoInit(&mr, pVnode->pMeta, META_READER_LOCK, 0);
39,044 ✔
848
  if (metaGetTableEntryByName(&mr, tableName) != 0) {
39,044 ✔
849
    vDebug("vgId:%d %s ref table not exist: %s", TD_VID(pVnode), __func__, tableName);
×
850
    // Fill all columns with table-not-exist error
851
    for (int32_t i = 0; i < nCols; ++i) {
×
852
      SVTableRefResolveRspItem r = {0};
×
853
      r.code = TSDB_CODE_STREAM_VTB_REF_TABLE_NOT_EXIST;
×
854
      if (taosArrayPush(pRspItems, &r) == NULL) { code = terrno; break; }
×
855
    }
856
    metaReaderClear(&mr);
×
857
    return code;
×
858
  }
859

860
  bool isVtable = (mr.me.type == TSDB_VIRTUAL_NORMAL_TABLE || mr.me.type == TSDB_VIRTUAL_CHILD_TABLE);
39,044 ✔
861

862
  // Release mr's meta read lock before opening any further LOCK readers below
863
  // (stbReader, and the per-column tag-value readers inside vnodeFillResolveRspFromColRef):
864
  // a nested META_READER_LOCK rdlock deadlocks once a writer is queued on the
865
  // meta rwlock (glibc blocks new readers behind a pending writer).
866
  // mr.me stays valid until metaReaderClear.
867
  metaReaderReleaseLock(&mr);
39,044 ✔
868

869
  // Pre-read parent stable info for virtual child table (shared across all columns)
870
  SMetaReader stbReader      = {0};
39,044 ✔
871
  bool        stbReaderInited = false;
39,044 ✔
872
  if (isVtable && mr.me.type == TSDB_VIRTUAL_CHILD_TABLE) {
39,044 ✔
873
    metaReaderDoInit(&stbReader, pVnode->pMeta, META_READER_LOCK, 0);
3,015 ✔
874
    if (metaReaderGetTableEntryByUid(&stbReader, mr.me.ctbEntry.suid) != 0) {
3,015 ✔
875
      // Parent stable not found: all columns fail
876
      for (int32_t i = 0; i < nCols; ++i) {
×
877
        SVTableRefResolveRspItem r = {0};
×
878
        r.code = TSDB_CODE_STREAM_VTB_REF_TABLE_NOT_EXIST;
×
879
        if (taosArrayPush(pRspItems, &r) == NULL) { code = terrno; break; }
×
880
      }
881
      metaReaderClear(&stbReader);
×
882
      metaReaderClear(&mr);
×
883
      return code;
×
884
    }
885
    stbReaderInited = true;
3,015 ✔
886
    // Same as above: drop the lock before the per-column loop, which opens
887
    // further LOCK readers via vnodeFillResolveRspFromColRef. stbReader.me
888
    // stays valid until metaReaderClear.
889
    metaReaderReleaseLock(&stbReader);
3,015 ✔
890
  }
891

892
  for (int32_t ci = 0; ci < nCols; ++ci) {
91,673 ✔
893
    SVTableRefResolveColSpec *c = taosArrayGet(pCols, ci);
52,629 ✔
894
    SVTableRefResolveRspItem  r = {0};
52,629 ✔
895

896
    SColRef *pFound = NULL;
52,629 ✔
897
    if (isVtable) {
52,629 ✔
898
      (void)vnodeFindVTableColRef(pVnode, &mr.me, stbReaderInited ? &stbReader.me : NULL,
7,437 ✔
899
                                  c->kind, c->colName, &pFound);
7,437 ✔
900
    }
901
    (void)vnodeFillResolveRspFromColRef(pVnode, &mr.me, pFound, c->kind,
52,629 ✔
902
                                        dbName, tableName, c->colName, isVtable, &r);
52,629 ✔
903

904
    if (taosArrayPush(pRspItems, &r) == NULL) {
52,629 ✔
905
      taosMemoryFreeClear(r.tagData);
×
906
      code = terrno;
×
907
      break;
×
908
    }
909
  }
910

911
  if (stbReaderInited) metaReaderClear(&stbReader);
39,044 ✔
912
  metaReaderClear(&mr);
39,044 ✔
913
  return code;
39,044 ✔
914
}
915

916
static int32_t vnodeResolveOneHop(SVnode *pVnode, const SVTableRefResolveItem *q,
46,089 ✔
917
                                  SVTableRefResolveRspItem *r) {
918
  SMetaReader mr   = {0};
46,089 ✔
919
  int32_t     code = 0;
46,089 ✔
920

921
  vDebug("vgId:%d %s enter: kind=%d ref=%s.%s.%s", TD_VID(pVnode), __func__, q->kind, q->refDbName,
46,089 ✔
922
         q->refTableName, q->refColName);
923

924
  metaReaderDoInit(&mr, pVnode->pMeta, META_READER_LOCK, 0);
46,089 ✔
925
  if (metaGetTableEntryByName(&mr, q->refTableName) != 0) {
46,089 ✔
926
    vDebug("vgId:%d %s ref table not exist: %s", TD_VID(pVnode), __func__, q->refTableName);
67 ✔
927
    r->code = TSDB_CODE_STREAM_VTB_REF_TABLE_NOT_EXIST;
67 ✔
928
    metaReaderClear(&mr);
67 ✔
929
    return 0;
67 ✔
930
  }
931

932
  bool isVtable = (mr.me.type == TSDB_VIRTUAL_NORMAL_TABLE || mr.me.type == TSDB_VIRTUAL_CHILD_TABLE);
46,022 ✔
933
  vDebug("vgId:%d %s table found: name=%s type=%d isVtable=%d", TD_VID(pVnode), __func__,
46,022 ✔
934
         q->refTableName, mr.me.type, isVtable);
935

936
  // Release mr's meta read lock before the nested LOCK readers opened by
937
  // vnodeFindVTableColRef (tmpStb) and vnodeFillResolveRspFromColRef
938
  // (streamReadChildTagConstValueImpl): a nested rdlock deadlocks once a writer
939
  // is queued on the meta rwlock. mr.me stays valid until metaReaderClear.
940
  metaReaderReleaseLock(&mr);
46,022 ✔
941

942
  if (isVtable) {
46,022 ✔
943
    // Lookup the SColRef via shared helper. Pass NULL for pStbEntry — for
944
    // single-hop the per-call parent-stable open cost is acceptable.
945
    SColRef *pFound = NULL;
16,281 ✔
946
    int32_t  rc     = vnodeFindVTableColRef(pVnode, &mr.me, NULL, q->kind, q->refColName, &pFound);
16,281 ✔
947
    if (rc != 0) {
16,281 ✔
948
      r->code = rc;
×
949
      metaReaderClear(&mr);
×
950
      return 0;
×
951
    }
952

953
    (void)vnodeFillResolveRspFromColRef(pVnode, &mr.me, pFound, q->kind,
16,281 ✔
954
                                        q->refDbName, q->refTableName, q->refColName,
16,281 ✔
955
                                        true, r);
956
    metaReaderClear(&mr);
16,281 ✔
957
    return 0;
16,281 ✔
958
  }
959

960
  // Physical table
961
  (void)vnodeFillResolveRspFromColRef(pVnode, &mr.me, NULL, q->kind,
29,741 ✔
962
                                      q->refDbName, q->refTableName, q->refColName,
29,741 ✔
963
                                      false, r);
964
  metaReaderClear(&mr);
29,741 ✔
965
  return 0;
29,741 ✔
966
}
967

968
int32_t vnodeProcessVTableRefResolveReq(SVnode *pVnode, SRpcMsg *pMsg) {
34,136 ✔
969
  int32_t              code   = 0;
34,136 ✔
970
  int32_t              rspLen = 0;
34,136 ✔
971
  void                *pBuf   = NULL;
34,136 ✔
972
  SVTableRefResolveReq req    = {0};
34,136 ✔
973
  SVTableRefResolveRsp rsp    = {0};
34,076 ✔
974
  SRpcMsg              rspMsg = {0};
34,076 ✔
975

976
  vTrace("vgId:%d %s enter: contLen=%d msgType=%d", TD_VID(pVnode), __func__, pMsg->contLen,
34,076 ✔
977
         pMsg->msgType);
978

979
  if (tDeserializeSVTableRefResolveReq((char *)pMsg->pCont + sizeof(SMsgHead),
34,136 ✔
980
                                       pMsg->contLen - (int32_t)sizeof(SMsgHead), &req) < 0) {
34,076 ✔
981
    vError("vgId:%d %s deserialize failed", TD_VID(pVnode), __func__);
×
982
    code = TSDB_CODE_INVALID_MSG;
×
983
    goto _end;
×
984
  }
985

986
  {
987
    // Table-grouped format: resolve per-table batch (meta opened once per table)
988
    int32_t nGroups = (req.groups != NULL) ? (int32_t)taosArrayGetSize(req.groups) : 0;
34,136 ✔
989
    // Count total columns across all groups for pre-allocation
990
    int32_t totalCols = 0;
34,136 ✔
991
    for (int32_t i = 0; i < nGroups; ++i) {
73,180 ✔
992
      SVTableRefResolveGroupItem *g = taosArrayGet(req.groups, i);
39,044 ✔
993
      totalCols += (g->cols != NULL) ? (int32_t)taosArrayGetSize(g->cols) : 0;
39,044 ✔
994
    }
995
    vTrace("vgId:%d %s req: ver=%" PRId64 " groups=%d totalCols=%d",
34,136 ✔
996
           TD_VID(pVnode), __func__, req.ver, nGroups, totalCols);
997

998
    rsp.items = taosArrayInit(totalCols, sizeof(SVTableRefResolveRspItem));
34,136 ✔
999
    if (rsp.items == NULL) { code = terrno; goto _end; }
34,136 ✔
1000

1001
    for (int32_t i = 0; i < nGroups; ++i) {
73,180 ✔
1002
      SVTableRefResolveGroupItem *g = taosArrayGet(req.groups, i);
39,044 ✔
1003
      int32_t rc = vnodeResolveTableGroup(pVnode, g->dbName, g->tableName, g->cols, rsp.items);
39,044 ✔
1004
      if (rc != 0) {
39,044 ✔
1005
        code = rc;
×
1006
        goto _end;
×
1007
      }
1008
    }
1009
  }
1010

1011
  rspLen = tSerializeSVTableRefResolveRsp(NULL, 0, &rsp);
34,136 ✔
1012
  if (rspLen < 0) {
34,072 ✔
1013
    code = TSDB_CODE_OUT_OF_MEMORY;
×
1014
    goto _end;
×
1015
  }
1016
  pBuf = rpcMallocCont(rspLen);
34,072 ✔
1017
  if (pBuf == NULL) {
34,136 ✔
1018
    code = terrno;
×
1019
    goto _end;
×
1020
  }
1021
  if (tSerializeSVTableRefResolveRsp(pBuf, rspLen, &rsp) < 0) {
34,136 ✔
1022
    rpcFreeCont(pBuf);
×
1023
    pBuf = NULL;
×
1024
    code = TSDB_CODE_OUT_OF_MEMORY;
×
1025
    goto _end;
×
1026
  }
1027

1028
_end:
34,136 ✔
1029
  tFreeSVTableRefResolveReq(&req);
34,136 ✔
1030
  tFreeSVTableRefResolveRsp(&rsp);
34,072 ✔
1031

1032
  rspMsg.info    = pMsg->info;
34,076 ✔
1033
  rspMsg.pCont   = (code == 0) ? pBuf : NULL;
34,076 ✔
1034
  rspMsg.contLen = (code == 0) ? rspLen : 0;
34,076 ✔
1035
  rspMsg.code    = code;
34,076 ✔
1036
  rspMsg.msgType = pMsg->msgType;
34,076 ✔
1037

1038
  if (code != 0) {
34,076 ✔
1039
    vError("vgId:%d, vtable ref resolve failed since %s", TD_VID(pVnode), tstrerror(code));
×
1040
    if (pBuf != NULL) rpcFreeCont(pBuf);
×
1041
  }
1042

1043
  vDebug("vgId:%d %s send rsp: code=0x%x rspLen=%d", TD_VID(pVnode), __func__, code, rspMsg.contLen);
34,076 ✔
1044
  tmsgSendRsp(&rspMsg);
34,076 ✔
1045
  return 0;
34,136 ✔
1046
}
1047

1048
// ============================================================================
1049
// C2: multi-hop chain driver (client side)
1050
// Task 6: chain resolution loop (single-vgId / local-vnode simplified version)
1051
// ============================================================================
1052

1053
#define STREAM_VTB_MAX_HOPS 32
1054

1055
static void freeColMap(void* ptr) {
136,465 ✔
1056
  taosMemoryFree(*(void**)ptr); 
136,465 ✔
1057
}
136,465 ✔
1058

1059
static void freeTagMap(void* ptr) {
28,725 ✔
1060
  STagValue **pp = (STagValue **)ptr;
28,725 ✔
1061
  if (*pp) taosMemoryFreeClear((*pp)->pData);
28,725 ✔
1062
  taosMemoryFreeClear(*pp);
28,725 ✔
1063
}
28,725 ✔
1064

1065
// (SResolveWorkItem moved to vnodeStreamVTable.h for testability.)
1066
static SVTableResolveResult *streamGetOrCreateUidResult(SSHashObj *uid2Result, int64_t uid) {
171,721 ✔
1067
  SVTableResolveResult **ppRes = (SVTableResolveResult **)tSimpleHashGet(uid2Result, &uid, sizeof(uid));
171,721 ✔
1068
  if (ppRes != NULL && *ppRes != NULL) {
171,721 ✔
1069
    return *ppRes;
103,363 ✔
1070
  }
1071

1072
  SVTableResolveResult *pRes = taosMemoryCalloc(1, sizeof(*pRes));
68,358 ✔
1073
  if (pRes == NULL) return NULL;
68,291 ✔
1074
  pRes->colMap = tSimpleHashInit(8, taosGetDefaultHashFunction(TSDB_DATA_TYPE_SMALLINT));
68,291 ✔
1075
  pRes->tagMap = tSimpleHashInit(8, taosGetDefaultHashFunction(TSDB_DATA_TYPE_SMALLINT));
68,358 ✔
1076

1077
  if (pRes->colMap == NULL || pRes->tagMap == NULL) {
68,358 ✔
1078
    streamVTableResolveResultDestroy(&pRes);
×
1079
    return NULL;
×
1080
  }
1081
  tSimpleHashSetFreeFp(pRes->colMap, freeColMap);
68,358 ✔
1082
  tSimpleHashSetFreeFp(pRes->tagMap, freeTagMap);
68,358 ✔
1083

1084
  if (tSimpleHashPut(uid2Result, &uid, sizeof(uid), &pRes, sizeof(pRes)) != 0) {
68,358 ✔
1085
    streamVTableResolveResultDestroy(&pRes);
×
1086
    return NULL;
×
1087
  }
1088
  return pRes;
68,299 ✔
1089
}
1090

1091
// Push initial work-items for a single vtable uid. Each requested cid (col or tag)
1092
// is resolved against the local vtable entry's pColRef / pTagRef:
1093
//   - COL hasRef=true   -> push next-hop work-item
1094
//   - COL hasRef=false  -> directly write terminal SColResolveItem{hasRef=false} into colMap
1095
//   - TAG hasRef=true   -> push next-hop work-item
1096
//   - TAG hasRef=false  -> on a virtual child table the tag may be stored as a
1097
//                          constant value on the vchild's own STag; read it locally
1098
//                          and write a terminal STagValue into tagMap (no work-item).
1099
//                          Virtual normal tables have no tag concept and fail.
1100
//
1101
// colCids == NULL means "all columns of this vtable" (used by the only-ts trigger
1102
// path where the request carries just the primary-key TS but the response must
1103
// describe every column ref). tagCids == NULL is treated as "no tag".
1104
// Returns 0 on success; non-zero means whole-uid skip (table missing / cid missing / OOM).
1105
static int32_t streamPushInitialWorkItemsForUid(SVnode *pVnode, int64_t uid, SArray *colCids, SArray *tagCids,
77,904 ✔
1106
                                                SArray *workList, SSHashObj *uid2Result) {
1107
  int32_t     code = 0;
77,904 ✔
1108
  SMetaReader mr   = {0};
77,904 ✔
1109
  metaReaderDoInit(&mr, pVnode->pMeta, META_READER_LOCK, 0);
77,904 ✔
1110

1111
  // H2 v0.5: top-level vtable uid not present in local meta (concurrently
1112
  // dropped) or entry type is not a vtable. Treat as a soft skip: log a
1113
  // warning and return 0 without producing any uid2Result entry. The caller
1114
  // (streamResolveVTableRefChain seed loop) sees rc==0 and simply continues;
1115
  // downstream consumers that strictly require this uid (e.g. PSEUDO_COL
1116
  // single-uid path) detect the missing entry and raise the error.
1117
  if (metaReaderGetTableEntryByUid(&mr, uid) != 0) {
77,904 ✔
1118
    stWarn("vgId:%d %s uid=%" PRId64 " META_NOT_FOUND -> H2 skip", TD_VID(pVnode), __func__, uid);
222 ✔
1119
    goto _end;
222 ✔
1120
  }
1121
  if (mr.me.type != TSDB_VIRTUAL_NORMAL_TABLE && mr.me.type != TSDB_VIRTUAL_CHILD_TABLE) {
77,682 ✔
1122
    stWarn("vgId:%d %s uid=%" PRId64 " type=%d not vtable -> H2 skip",
×
1123
           TD_VID(pVnode), __func__, uid, mr.me.type);
1124
    goto _end;
×
1125
  }
1126

1127
  // Release mr's meta read lock: the tag loop below opens a nested LOCK reader
1128
  // via streamReadChildTagConstValueByCid, and a nested rdlock deadlocks once a
1129
  // writer is queued on the meta rwlock. mr.me stays valid until metaReaderClear.
1130
  metaReaderReleaseLock(&mr);
77,682 ✔
1131

1132
  SVTableResolveResult *pRes = streamGetOrCreateUidResult(uid2Result, uid);
77,682 ✔
1133
  if (pRes == NULL) {
77,615 ✔
1134
    code = terrno;
×
1135
    goto _end;
×
1136
  }
1137

1138
  // Resolve column cids against pColRef. When colCids==NULL, iterate every
1139
  // entry of this vtable's pColRef directly (no per-cid lookup needed).
1140
  if (colCids == NULL) {
77,615 ✔
1141
    for (int32_t j = 0; j < mr.me.colRef.nCols; ++j) {
91,813 ✔
1142
      SColRef *pRef = &mr.me.colRef.pColRef[j];
64,810 ✔
1143
      col_id_t cid  = pRef->id;
64,810 ✔
1144
      if (!pRef->hasRef) {
64,810 ✔
1145
        SColResolveItem *item = taosMemoryCalloc(1, sizeof(*item));
26,944 ✔
1146
        if (item == NULL) { code = terrno; goto _end; }
26,944 ✔
1147
        item->hasRef = false;
26,944 ✔
1148
        // Snapshot old pointer before put; free it only after successful put so
1149
        // the hash never holds a dangling pointer (avoids double-free on cleanup).
1150
        SColResolveItem **ppOld = (SColResolveItem **)tSimpleHashGet(pRes->colMap, &cid, sizeof(cid));
26,944 ✔
1151
        SColResolveItem  *oldItem = (ppOld && *ppOld) ? *ppOld : NULL;
27,003 ✔
1152
        if (tSimpleHashPut(pRes->colMap, &cid, sizeof(cid), &item, sizeof(item)) != 0) {
27,003 ✔
1153
          taosMemoryFree(item);
×
1154
          code = terrno;
×
1155
          goto _end;
×
1156
        }
1157
        if (oldItem) { taosMemoryFree(oldItem); }
26,885 ✔
1158
        continue;
26,885 ✔
1159
      }
1160
      SResolveWorkItem w = {0};
37,866 ✔
1161
      w.originVtbUid = uid;
37,866 ✔
1162
      w.originCid    = cid;
37,866 ✔
1163
      w.kind         = STREAM_VREF_KIND_COL;
37,866 ✔
1164
      tstrncpy(w.refDbName,    pRef->refDbName,    TSDB_DB_NAME_LEN);
37,866 ✔
1165
      tstrncpy(w.refTableName, pRef->refTableName, TSDB_TABLE_NAME_LEN);
37,866 ✔
1166
      tstrncpy(w.refColName,   pRef->refColName,   TSDB_COL_NAME_LEN);
37,866 ✔
1167
      if (taosArrayPush(workList, &w) == NULL) { code = terrno; goto _end; }
37,807 ✔
1168
    }
1169
  } else {
1170
    int32_t nCol = (int32_t)taosArrayGetSize(colCids);
50,612 ✔
1171
    for (int32_t i = 0; i < nCol; ++i) {
140,998 ✔
1172
      col_id_t cid    = *(col_id_t *)taosArrayGet(colCids, i);
90,319 ✔
1173
      SColRef *pFound = NULL;
90,319 ✔
1174
      for (int32_t j = 0; j < mr.me.colRef.nCols; ++j) {
147,706 ✔
1175
        if (mr.me.colRef.pColRef[j].id == cid) {
147,706 ✔
1176
          pFound = &mr.me.colRef.pColRef[j];
90,319 ✔
1177
          break;
90,319 ✔
1178
        }
1179
      }
1180
      if (pFound == NULL) {
90,319 ✔
1181
        stWarn("vgId:%d %s uid=%" PRId64 " COL cid=%d NOT_IN_COLREF -> uid skip",
×
1182
               TD_VID(pVnode), __func__, uid, cid);
1183
        code = TSDB_CODE_STREAM_VTB_REF_COL_NOT_EXIST;
×
1184
        goto _end;
×
1185
      }
1186
      if (!pFound->hasRef) {
90,319 ✔
1187
        SColResolveItem *item = taosMemoryCalloc(1, sizeof(*item));
47,203 ✔
1188
        if (item == NULL) { code = terrno; goto _end; }
47,203 ✔
1189
        item->hasRef = false;
47,203 ✔
1190
        if (tSimpleHashPut(pRes->colMap, &cid, sizeof(cid), &item, sizeof(item)) != 0) {
47,203 ✔
1191
          taosMemoryFree(item);
×
1192
          code = terrno;
×
1193
          goto _end;
×
1194
        }
1195
        continue;
47,203 ✔
1196
      }
1197
      SResolveWorkItem w = {0};
43,116 ✔
1198
      w.originVtbUid = uid;
43,116 ✔
1199
      w.originCid    = cid;
43,116 ✔
1200
      w.kind         = STREAM_VREF_KIND_COL;
43,116 ✔
1201
      tstrncpy(w.refDbName,    pFound->refDbName,    TSDB_DB_NAME_LEN);
43,116 ✔
1202
      tstrncpy(w.refTableName, pFound->refTableName, TSDB_TABLE_NAME_LEN);
43,116 ✔
1203
      tstrncpy(w.refColName,   pFound->refColName,   TSDB_COL_NAME_LEN);
43,116 ✔
1204
      if (taosArrayPush(workList, &w) == NULL) { code = terrno; goto _end; }
43,116 ✔
1205
    }
1206
  }
1207

1208
  // resolve tag cids against pTagRef
1209
  int32_t nTag = (tagCids != NULL) ? (int32_t)taosArrayGetSize(tagCids) : 0;
77,682 ✔
1210
  for (int32_t i = 0; i < nTag; ++i) {
106,281 ✔
1211
    col_id_t cid    = *(col_id_t *)taosArrayGet(tagCids, i);
28,540 ✔
1212
    SColRef *pFound = NULL;
28,725 ✔
1213
    for (int32_t j = 0; j < mr.me.colRef.nTagRefs; ++j) {
30,467 ✔
1214
      if (mr.me.colRef.pTagRef[j].id == cid) {
14,807 ✔
1215
        pFound = &mr.me.colRef.pTagRef[j];
13,065 ✔
1216
        break;
13,065 ✔
1217
      }
1218
    }
1219
    // For a VCT, a tag cid that is absent from pTagRef[] (or present with
1220
    // hasRef=0) means the value is stored locally on the child entry as a
1221
    // constant inherited from the parent vstable schemaTag. Both cases must
1222
    // go through the local constant-read path. Only non-VCT (i.e. VNT) tags
1223
    // truly do not exist and should skip the uid.
1224
    if (pFound == NULL && mr.me.type != TSDB_VIRTUAL_CHILD_TABLE) {
28,725 ✔
1225
      stWarn("vgId:%d %s uid=%" PRId64 " TAG cid=%d NOT_IN_TAGREF type=%d -> uid skip",
×
1226
             TD_VID(pVnode), __func__, uid, cid, mr.me.type);
1227
      code = TSDB_CODE_STREAM_VTB_REF_COL_NOT_EXIST;
×
1228
      goto _end;
×
1229
    }
1230

1231
    if (pFound == NULL || !pFound->hasRef) {
28,725 ✔
1232
      // Constant tag on a virtual child table: read locally, write terminal STagValue.
1233
      STagValue *tv = taosMemoryCalloc(1, sizeof(*tv));
15,593 ✔
1234
      if (tv == NULL) { code = terrno; goto _end; }
15,593 ✔
1235
      // Use cid: SColRef.colName is not persisted for vtable on disk
1236
      // (the field is "for tmq get json" only). Resolve tag by colId in stable schemaTag.
1237
      int32_t rc = streamReadChildTagConstValueByCid(pVnode, &mr.me, cid,
15,593 ✔
1238
                                                    &tv->type, &tv->nLen, &tv->pData);
15,593 ✔
1239
      if (rc != 0) {
15,534 ✔
1240
        stWarn("vgId:%d %s uid=%" PRId64 " TAG cid=%d const-read err=0x%x -> uid skip",
×
1241
               TD_VID(pVnode), __func__, uid, cid, rc);
1242
        taosMemoryFreeClear(tv->pData);
×
1243
        taosMemoryFree(tv);
×
1244
        code = rc;
×
1245
        goto _end;
×
1246
      }
1247
      if (tSimpleHashPut(pRes->tagMap, &cid, sizeof(cid), &tv, sizeof(tv)) != 0) {
15,534 ✔
1248
        taosMemoryFreeClear(tv->pData);
×
1249
        taosMemoryFree(tv);
×
1250
        code = terrno;
×
1251
        goto _end;
×
1252
      }
1253
      continue;
15,475 ✔
1254
    }
1255

1256
    SResolveWorkItem w = {0};
13,132 ✔
1257
    w.originVtbUid = uid;
13,065 ✔
1258
    w.originCid    = cid;
13,065 ✔
1259
    w.kind         = STREAM_VREF_KIND_TAG;
13,065 ✔
1260
    tstrncpy(w.refDbName,    pFound->refDbName,    TSDB_DB_NAME_LEN);
13,065 ✔
1261
    tstrncpy(w.refTableName, pFound->refTableName, TSDB_TABLE_NAME_LEN);
13,065 ✔
1262
    tstrncpy(w.refColName,   pFound->refColName,   TSDB_COL_NAME_LEN);
13,065 ✔
1263
    if (taosArrayPush(workList, &w) == NULL) { code = terrno; goto _end; }
13,065 ✔
1264
  }
1265

1266
_end:
77,904 ✔
1267
  metaReaderClear(&mr);
77,837 ✔
1268
  return code;
77,904 ✔
1269
}
1270

1271
// Local hash comparator: search a hash value in a sorted SArray<SVgroupInfo>.
1272
// Mirrors catalog/ctgUtil.c:ctgHashValueComp; keep them in sync.
1273
static int32_t streamVgHashValueComp(void const *lp, void const *rp) {
139,847 ✔
1274
  uint32_t    *key = (uint32_t *)lp;
139,847 ✔
1275
  SVgroupInfo *pVg = (SVgroupInfo *)rp;
139,847 ✔
1276
  if (*key < pVg->hashBegin) return -1;
139,847 ✔
1277
  if (*key > pVg->hashEnd)   return 1;
120,370 ✔
1278
  return 0;
98,718 ✔
1279
}
1280

1281
static int32_t streamVgInfoBeginComp(void const *lp, void const *rp) {
55,051 ✔
1282
  SVgroupInfo *pLeft  = (SVgroupInfo *)lp;
55,051 ✔
1283
  SVgroupInfo *pRight = (SVgroupInfo *)rp;
55,051 ✔
1284
  if (pLeft->hashBegin < pRight->hashBegin) return -1;
55,051 ✔
1285
  if (pLeft->hashBegin > pRight->hashBegin) return 1;
×
1286
  return 0;
×
1287
}
1288

1289
// Async-callback context used by streamFetchDbVgInfo to receive SUseDbRsp.
1290
// Heap-allocated so the callback can safely access it even after caller timeout.
1291
//
1292
// Ownership is transferred via the `state` CAS:
1293
//   FETCH_DBVG_INFLIGHT(0) -> FETCH_DBVG_CB_DONE(1)    : callback finished, driver owns ctx
1294
//   FETCH_DBVG_INFLIGHT(0) -> FETCH_DBVG_DRIVER_GONE(2): driver timed out, callback owns ctx
1295
// Whichever side wins the CAS is responsible for releasing ctx (and destroying the sem).
1296
#define FETCH_DBVG_INFLIGHT    0
1297
#define FETCH_DBVG_CB_DONE     1
1298
#define FETCH_DBVG_DRIVER_GONE 2
1299

1300
typedef struct SStreamFetchDbVgCtx {
1301
  tsem2_t    ready;
1302
  SUseDbRsp *pRsp;
1303
  int32_t    code;
1304
  int8_t     state;
1305
} SStreamFetchDbVgCtx;
1306

1307
static void streamDestroyFetchDbVgCtx(SStreamFetchDbVgCtx *pCtx) {
17,743 ✔
1308
  if (pCtx == NULL) return;
17,743 ✔
1309
  if (pCtx->pRsp != NULL) {
17,743 ✔
1310
    tFreeSUsedbRsp(pCtx->pRsp);
×
1311
    taosMemoryFree(pCtx->pRsp);
×
1312
  }
1313
  TAOS_UNUSED(tsem2_destroy(&pCtx->ready));
17,743 ✔
1314
  taosMemoryFree(pCtx);
17,743 ✔
1315
}
1316

1317
static int32_t streamProcessFetchDbVgRsp(void *param, SDataBuf *pMsg, int32_t code) {
17,743 ✔
1318
  SStreamFetchDbVgCtx *pCtx = (SStreamFetchDbVgCtx *)param;
17,743 ✔
1319
  if (code == TSDB_CODE_SUCCESS && pMsg != NULL && pMsg->pData != NULL && pMsg->len > 0) {
17,743 ✔
1320
    pCtx->pRsp = taosMemoryCalloc(1, sizeof(SUseDbRsp));
17,743 ✔
1321
    if (pCtx->pRsp == NULL) {
17,743 ✔
1322
      code = terrno;
×
1323
    } else if (tDeserializeSUseDbRsp(pMsg->pData, (int32_t)pMsg->len, pCtx->pRsp) != 0) {
17,743 ✔
1324
      code = TSDB_CODE_INVALID_MSG;
×
1325
    }
1326
  } else if (code == TSDB_CODE_SUCCESS) {
×
1327
    code = TSDB_CODE_INVALID_MSG;
×
1328
  }
1329
  pCtx->code = code;
17,743 ✔
1330

1331
  if (pMsg != NULL) {
17,743 ✔
1332
    taosMemoryFreeClear(pMsg->pData);
17,743 ✔
1333
    taosMemoryFreeClear(pMsg->pEpSet);
17,743 ✔
1334
  }
1335

1336
  // Atomically claim ownership: CAS INFLIGHT -> CB_DONE.
1337
  // If CAS fails, the driver already abandoned ctx; callback must clean up.
1338
  int8_t prev = atomic_val_compare_exchange_8(&pCtx->state, FETCH_DBVG_INFLIGHT, FETCH_DBVG_CB_DONE);
17,743 ✔
1339
  if (prev == FETCH_DBVG_DRIVER_GONE) {
17,743 ✔
1340
    streamDestroyFetchDbVgCtx(pCtx);
×
1341
    return code;
×
1342
  }
1343

1344
  // Driver still owns ctx and is (or will be) waiting on the sem.
1345
  TAOS_UNUSED(tsem2_post(&pCtx->ready));
17,743 ✔
1346
  return code;
17,743 ✔
1347
}
1348

1349
// Fetch SUseDbRsp asynchronously from mnode for dbFName ("acctId.dbName").
1350
// Caller owns *ppOut on success and must call tFreeSUsedbRsp + free the pointer.
1351
static int32_t streamFetchDbVgInfo(SVnode *pVnode, const char *dbFName, SUseDbRsp **ppOut) {
17,743 ✔
1352
  int32_t              code      = 0;
17,743 ✔
1353
  SUseDbReq            req       = {0};
17,743 ✔
1354
  void                *pReqBuf   = NULL;
17,743 ✔
1355
  SMsgSendInfo        *pSendInfo = NULL;
17,743 ✔
1356
  SStreamFetchDbVgCtx *pCtx      = NULL;
17,743 ✔
1357
  SEpSet               epSet     = {0};
17,743 ✔
1358

1359
  *ppOut = NULL;
17,743 ✔
1360
  tstrncpy(req.db, dbFName, sizeof(req.db));
17,743 ✔
1361
  req.vgVersion  = -1;
17,743 ✔
1362
  req.dbId       = 0;
17,743 ✔
1363
  req.numOfTable = 0;
17,743 ✔
1364
  req.stateTs    = 0;
17,743 ✔
1365

1366
  void *clientRpc = pVnode->msgCb.clientRpc;
17,743 ✔
1367
  if (clientRpc == NULL) { code = TSDB_CODE_INVALID_PARA; goto _end; }
17,743 ✔
1368

1369
  // Heap-allocate ctx so callback can safely access it after caller timeout.
1370
  pCtx = taosMemoryCalloc(1, sizeof(SStreamFetchDbVgCtx));
17,743 ✔
1371
  if (pCtx == NULL) { code = terrno; goto _end; }
17,743 ✔
1372
  if (tsem2_init(&pCtx->ready, 0, 0) != 0) {
17,743 ✔
1373
    // sem not initialized yet; free directly to avoid destroying an uninitialized sem.
1374
    code = terrno;
×
1375
    taosMemoryFree(pCtx);
×
1376
    pCtx = NULL;
×
1377
    goto _end;
×
1378
  }
1379

1380
  int32_t reqLen = tSerializeSUseDbReq(NULL, 0, &req);
17,743 ✔
1381
  if (reqLen < 0) { code = terrno; goto _end; }
17,743 ✔
1382
  pReqBuf = taosMemoryCalloc(1, reqLen);
17,743 ✔
1383
  if (pReqBuf == NULL) { code = terrno; goto _end; }
17,684 ✔
1384
  if (tSerializeSUseDbReq(pReqBuf, reqLen, &req) < 0) { code = terrno; goto _end; }
17,684 ✔
1385

1386
  pSendInfo = taosMemoryCalloc(1, sizeof(SMsgSendInfo));
17,743 ✔
1387
  if (pSendInfo == NULL) { code = terrno; goto _end; }
17,684 ✔
1388

1389
  pSendInfo->param          = pCtx;
17,684 ✔
1390
  pSendInfo->msgInfo.pData  = pReqBuf;
17,684 ✔
1391
  pSendInfo->msgInfo.len    = reqLen;
17,684 ✔
1392
  pSendInfo->msgType        = TDMT_MND_GET_DB_INFO;
17,684 ✔
1393
  pSendInfo->fp             = streamProcessFetchDbVgRsp;
17,684 ✔
1394
  pReqBuf = NULL;  // ownership transferred to pSendInfo
17,684 ✔
1395

1396
  streamGetMnodeEpset(&epSet);
17,684 ✔
1397

1398
  code = asyncSendMsgToServer(clientRpc, &epSet, NULL, pSendInfo);
17,743 ✔
1399
  pSendInfo = NULL;  // ownership transferred (freed by asyncSendMsgToServer on any path)
17,743 ✔
1400
  if (code != 0) goto _end;
17,743 ✔
1401

1402
  if (tsem2_timewait(&pCtx->ready, STREAM_VTB_RPC_TIMEOUT_MS) != 0) {
17,743 ✔
1403
    // Timewait reported timeout. Try to atomically claim "driver gives up".
1404
    // If CAS succeeds, callback (whenever it fires) will free ctx.
1405
    // If CAS fails, callback finished concurrently with the timeout; treat as success
1406
    // and fall through to the normal data-handling path. The matching post may or may
1407
    // not have already been observed -- drain it via tsem2_timewait so we can safely destroy.
1408
    int8_t prev = atomic_val_compare_exchange_8(&pCtx->state, FETCH_DBVG_INFLIGHT, FETCH_DBVG_DRIVER_GONE);
×
1409
    if (prev == FETCH_DBVG_INFLIGHT) {
×
1410
      stWarn("vgId:%d %s timeout waiting for mnode db-vg-info rsp for %s",
×
1411
             TD_VID(pVnode), __func__, dbFName);
1412
      pCtx = NULL;  // ownership transferred to callback
×
1413
      code = TSDB_CODE_TIMEOUT_ERROR;
×
1414
      goto _end;
×
1415
    }
1416
    // Late win by callback: consume the pending post (non-blocking by design).
1417
    TAOS_UNUSED(tsem2_wait(&pCtx->ready));
×
1418
    stDebug("vgId:%d %s timewait raced with callback; proceeding with received rsp for %s",
×
1419
            TD_VID(pVnode), __func__, dbFName);
1420
  }
1421

1422
  if (pCtx->code != 0) { code = pCtx->code; goto _end; }
17,743 ✔
1423
  if (pCtx->pRsp == NULL) { code = TSDB_CODE_INVALID_MSG; goto _end; }
17,743 ✔
1424

1425
  // Sort vgroup array by hashBegin so we can binary-search for routing.
1426
  if (pCtx->pRsp->pVgroupInfos != NULL) {
17,743 ✔
1427
    taosArraySort(pCtx->pRsp->pVgroupInfos, streamVgInfoBeginComp);
17,743 ✔
1428
  }
1429

1430
  *ppOut  = pCtx->pRsp;
17,743 ✔
1431
  pCtx->pRsp = NULL;
17,743 ✔
1432

1433
_end:
17,743 ✔
1434
  if (pReqBuf != NULL) taosMemoryFree(pReqBuf);
17,743 ✔
1435
  if (pSendInfo != NULL) taosMemoryFree(pSendInfo);
17,743 ✔
1436
  if (pCtx != NULL) {
17,743 ✔
1437
    // Either the RPC never went out (state still INFLIGHT, no callback will fire),
1438
    // or the callback already completed (state == CB_DONE). In both cases the driver
1439
    // owns pCtx and is responsible for releasing it. The handoff CAS in the timeout
1440
    // branch above sets pCtx=NULL on the abandoned path, so we never reach here with
1441
    // ownership transferred away.
1442
    streamDestroyFetchDbVgCtx(pCtx);
17,743 ✔
1443
  }
1444
  return code;
17,743 ✔
1445
}
1446

1447
// Get SUseDbRsp for dbFName, using cache if available; otherwise fetch and insert.
1448
// Returned *ppOut is owned by the cache (when pCache != NULL) or by the caller
1449
// (when pCache == NULL); caller never frees the cached entry.
1450
static int32_t streamGetOrFetchDbVgInfo(SVnode *pVnode, SStreamVTableInfoCache *pCache,
98,718 ✔
1451
                                        const char *dbFName, SUseDbRsp **ppOut, bool *pCached) {
1452
  *ppOut = NULL;
98,718 ✔
1453
  if (pCached) *pCached = false;
98,718 ✔
1454

1455
  if (pCache != NULL && pCache->dbVgInfo != NULL) {
98,718 ✔
1456
    SUseDbRsp *pHit = (SUseDbRsp *)taosHashGet(pCache->dbVgInfo, dbFName, strlen(dbFName));
84,542 ✔
1457
    if (pHit != NULL) {
84,542 ✔
1458
      *ppOut = pHit;
80,975 ✔
1459
      if (pCached) *pCached = true;
80,975 ✔
1460
      return 0;
80,975 ✔
1461
    }
1462
  }
1463

1464
  SUseDbRsp *pNew = NULL;
17,743 ✔
1465
  int32_t    code = streamFetchDbVgInfo(pVnode, dbFName, &pNew);
17,743 ✔
1466
  if (code != 0) return code;
17,743 ✔
1467

1468
  if (pCache != NULL && pCache->dbVgInfo != NULL) {
17,743 ✔
1469
    // taosHashPut copies the value bytes; we must still keep the inner array
1470
    // alive (pVgroupInfos is a heap pointer the cached entry now owns).
1471
    if (taosHashPut(pCache->dbVgInfo, dbFName, strlen(dbFName), pNew, sizeof(*pNew)) != 0) {
3,567 ✔
1472
      tFreeSUsedbRsp(pNew);
×
1473
      taosMemoryFree(pNew);
×
1474
      return terrno;
×
1475
    }
1476
    // Hash now owns pVgroupInfos via the copied struct; drop our outer wrapper
1477
    // without freeing the array (cleanup uses tFreeSUsedbRsp on hash entries).
1478
    taosMemoryFree(pNew);
3,567 ✔
1479
    *ppOut = (SUseDbRsp *)taosHashGet(pCache->dbVgInfo, dbFName, strlen(dbFName));
3,567 ✔
1480
    return 0;
3,567 ✔
1481
  }
1482

1483
  *ppOut = pNew;
14,176 ✔
1484
  return 0;
14,176 ✔
1485
}
1486

1487
// Resolve target vgId/epSet for a (db, table) using cached SUseDbRsp routing info.
1488
// dbFName is "acctId.dbName" (matches SUseDbRsp->db); tableName is the child name.
1489
static int32_t streamRouteTableToVg(SUseDbRsp *pRsp, const char *dbFName, const char *tableName,
98,718 ✔
1490
                                    int32_t *pVgId, SEpSet *pEpSet) {
1491
  if (pRsp == NULL || pRsp->pVgroupInfos == NULL) return TSDB_CODE_INVALID_PARA;
98,718 ✔
1492
  int32_t vgNum = (int32_t)taosArrayGetSize(pRsp->pVgroupInfos);
98,718 ✔
1493
  if (vgNum <= 0) return TSDB_CODE_MND_DB_NOT_EXIST;
98,718 ✔
1494

1495
  char fullName[TSDB_TABLE_FNAME_LEN] = {0};
98,718 ✔
1496
  int32_t n = tsnprintf(fullName, sizeof(fullName), "%s.%s", dbFName, tableName);
98,718 ✔
1497
  if (n <= 0) return TSDB_CODE_INVALID_PARA;
98,718 ✔
1498

1499
  uint32_t hashValue = (uint32_t)taosGetTbHashVal(fullName, n, pRsp->hashMethod,
98,718 ✔
1500
                                                  pRsp->hashPrefix, pRsp->hashSuffix);
98,718 ✔
1501
  SVgroupInfo *pVg = (SVgroupInfo *)taosArraySearch(pRsp->pVgroupInfos, &hashValue,
98,718 ✔
1502
                                                    streamVgHashValueComp, TD_EQ);
1503
  if (pVg == NULL) return TSDB_CODE_MND_DB_NOT_EXIST;
98,718 ✔
1504
  *pVgId  = pVg->vgId;
98,718 ✔
1505
  *pEpSet = pVg->epSet;
98,718 ✔
1506
  vDebug("stream route table:%s to vgId:%d, epSet inUse:%d numOfEps:%d",
98,718 ✔
1507
        fullName, pVg->vgId, pVg->epSet.inUse, pVg->epSet.numOfEps);
1508
  for (int32_t i = 0; i < pVg->epSet.numOfEps; ++i) {
197,436 ✔
1509
    vDebug("stream route table:%s vgId:%d ep[%d]: %s:%u",
98,718 ✔
1510
          fullName, pVg->vgId, i, pVg->epSet.eps[i].fqdn, pVg->epSet.eps[i].port);
1511
  }
1512
  return 0;
98,718 ✔
1513
}
1514

1515
// Shared fan-out completion sync owned by streamBatchFanoutDrain. All fired
1516
// handles point at the same instance via SStreamVgResolveCtx::pSync. The sync
1517
// is heap-allocated so it can outlive the drain frame on the timeout path:
1518
// any late callback can still safely deref pSync and (if it is the last
1519
// reference) free it.
1520
//
1521
// Two atomic counters:
1522
//   - pending: number of in-flight callbacks PLUS a "fan-out in progress"
1523
//     reservation held by the driver. Gates the binary sem. The cb that
1524
//     drives pending -> 0 posts the sem exactly once.
1525
//   - refs:    lifetime refcount. Driver holds 1, each successful fire holds 1.
1526
//     Whoever drives refs -> 0 (driver after wait, or the last cb on the
1527
//     timeout-abandoned path) calls streamFanoutSyncDestroy.
1528
//
1529
// Why two counters? The driver releases its `pending` reservation BEFORE
1530
// tsem2_timewait, but must keep `refs` so the sync stays alive while it's
1531
// blocked. On a real timeout the driver decrements `refs` and walks away;
1532
// the last late cb to dec refs frees the sync. This decouples "wakeup
1533
// signalling" from "object lifetime" and avoids a destroy-vs-late-post race.
1534
// SStreamFanoutSync moved to vnodeStreamVTable.h for testability.
1535

1536
SStreamFanoutSync *streamFanoutSyncCreate(void) {
45,632 ✔
1537
  SStreamFanoutSync *p = taosMemoryCalloc(1, sizeof(*p));
45,632 ✔
1538
  if (p == NULL) return NULL;
45,632 ✔
1539
  if (tsem2_init(&p->sem, 0, 0) != 0) {
45,632 ✔
1540
    taosMemoryFree(p);
×
1541
    return NULL;
×
1542
  }
1543
  atomic_store_32(&p->pending, 0);
45,632 ✔
1544
  atomic_store_32(&p->refs, 0);
45,632 ✔
1545
  return p;
45,632 ✔
1546
}
1547

1548
void streamFanoutSyncDestroy(SStreamFanoutSync *p) {
45,579 ✔
1549
  if (p == NULL) return;
45,579 ✔
1550
  TAOS_UNUSED(tsem2_destroy(&p->sem));
45,577 ✔
1551
  taosMemoryFree(p);
45,632 ✔
1552
}
1553

1554
// Release one ref to the shared sync. Caller MUST stop touching sync after
1555
// this returns. Returns true if this call freed the sync (last reference).
1556
bool streamFanoutSyncRelease(SStreamFanoutSync *p) {
79,715 ✔
1557
  if (p == NULL) return false;
79,715 ✔
1558
  int32_t r = atomic_sub_fetch_32(&p->refs, 1);
79,713 ✔
1559
  if (r == 0) {
79,768 ✔
1560
    streamFanoutSyncDestroy(p);
45,630 ✔
1561
    return true;
45,575 ✔
1562
  }
1563
  return false;
34,136 ✔
1564
}
1565

1566
// Async-callback context used by streamPrepareAndFireOneVgResolve /
1567
// streamScatterOneVgResolve to receive SVTableRefResolveRsp from a
1568
// remote vnode. One ctx is owned by each SStreamVgRpcHandle so callbacks can
1569
// fire any time during the fan-out window without racing handle teardown.
1570
typedef struct SStreamVgResolveCtx {
1571
  SStreamFanoutSync   *pSync;    // borrowed: shared completion sync owned by drain
1572
  SVTableRefResolveRsp rsp;
1573
  int32_t              code;
1574
} SStreamVgResolveCtx;
1575

1576
// Handle for an in-flight per-vg resolve RPC. Driver allocates one per remote
1577
// vg in fan-out phase A; phase B then waits + scatters + destroys them all,
1578
// possibly in parallel (each ctx is heap-resident so callbacks can fire any
1579
// time during phase A without racing with stack teardown).
1580
//
1581
// `state` is an atomic CAS-arbitrated ownership word used to hand off the
1582
// handle between callback and driver on the timeout path:
1583
//   VG_HANDLE_INFLIGHT(0)   : cb hasn't completed yet, driver owns
1584
//   VG_HANDLE_CB_DONE(1)    : cb completed normally, driver owns h
1585
//   VG_HANDLE_DRIVER_GONE(2): driver gave up on this handle, cb owns h
1586
// Whichever side wins the CAS becomes responsible for streamDestroyVgRpcHandle.
1587
#define VG_HANDLE_INFLIGHT    0
1588
#define VG_HANDLE_CB_DONE     1
1589
#define VG_HANDLE_DRIVER_GONE 2
1590

1591
typedef struct SStreamVgRpcHandle {
1592
  int32_t              vgId;
1593
  int8_t               state;           // atomic, see header comment
1594
  SArray              *indexList;       // borrowed: position list inside dedup batch
1595
  int32_t              totalCols;       // expected rsp.items count = sum of group cols
1596
  // scatterOrder[i] = dedupItems position of the i-th flattened column in req
1597
  // (groups[0].cols[0], groups[0].cols[1], ..., groups[1].cols[0], ...).
1598
  // This matches the server response order exactly and is the correct mapping
1599
  // to use during scatter, replacing the incorrect indexList-order assumption.
1600
  SArray              *scatterOrder;    // owned: SArray<int32_t>, freed on destroy
1601
  SVTableRefResolveReq req;             // owned: tFreeSVTableRefResolveReq on destroy
1602
  SStreamVgResolveCtx  ctx;            // owned: decoded rsp + shared-sem pointer
1603
} SStreamVgRpcHandle;
1604

1605
static void streamDestroyVgRpcHandle(SStreamVgRpcHandle **ppHandle) {
34,136 ✔
1606
  if (ppHandle == NULL || *ppHandle == NULL) return;
34,136 ✔
1607
  SStreamVgRpcHandle *h = *ppHandle;
34,136 ✔
1608
  tFreeSVTableRefResolveReq(&h->req);
34,136 ✔
1609
  tFreeSVTableRefResolveRsp(&h->ctx.rsp);
34,136 ✔
1610
  taosArrayDestroy(h->scatterOrder);
34,136 ✔
1611
  h->scatterOrder = NULL;
34,136 ✔
1612
  // Note: ctx.pSync is a refcounted shared sync; its lifetime is governed by
1613
  // streamFanoutSyncRelease (called by both driver and each cb) and is NOT
1614
  // released here.
1615
  taosMemoryFree(h);
34,136 ✔
1616
  *ppHandle = NULL;
34,081 ✔
1617
}
1618

1619
static int32_t streamProcessVgResolveRsp(void *param, SDataBuf *pMsg, int32_t code) {
34,136 ✔
1620
  SStreamVgResolveCtx *pCtx  = (SStreamVgResolveCtx *)param;
34,136 ✔
1621
  SStreamFanoutSync   *pSync = pCtx->pSync;
34,136 ✔
1622
  // Recover the enclosing handle from pCtx via offsetof: pCtx is the `ctx`
1623
  // member of the SStreamVgRpcHandle that owns it.
1624
  SStreamVgRpcHandle *h =
34,136 ✔
1625
      (SStreamVgRpcHandle *)((char *)pCtx - offsetof(SStreamVgRpcHandle, ctx));
34,072 ✔
1626

1627
  stTrace("stream vtable resolve rsp arrived: code=0x%x len=%d pData=%p", code,
34,072 ✔
1628
          pMsg ? (int32_t)pMsg->len : -1, pMsg ? pMsg->pData : NULL);
1629
  if (code == TSDB_CODE_SUCCESS) {
34,072 ✔
1630
    if (pMsg != NULL && pMsg->pData != NULL && pMsg->len > 0) {
34,072 ✔
1631
      if (tDeserializeSVTableRefResolveRsp(pMsg->pData, (int32_t)pMsg->len, &pCtx->rsp) < 0) {
34,072 ✔
1632
        code = TSDB_CODE_OUT_OF_MEMORY;
×
1633
      }
1634
    } else {
1635
      code = TSDB_CODE_INVALID_MSG;
×
1636
    }
1637
  }
1638
  pCtx->code = code;
34,136 ✔
1639
  stTrace("stream vtable resolve rsp processed: code=0x%x rspItems=%d", code,
34,136 ✔
1640
          pCtx->rsp.items ? (int32_t)taosArrayGetSize(pCtx->rsp.items) : 0);
1641

1642
  if (pMsg != NULL) {
34,136 ✔
1643
    taosMemoryFreeClear(pMsg->pData);
34,136 ✔
1644
    taosMemoryFreeClear(pMsg->pEpSet);
34,136 ✔
1645
  }
1646

1647
  // Claim ownership of the handle: CAS INFLIGHT -> CB_DONE. If CAS fails
1648
  // (state is VG_HANDLE_DRIVER_GONE), the driver abandoned this handle on
1649
  // the timeout path and the cb is responsible for releasing it.
1650
  int8_t prev = atomic_val_compare_exchange_8(&h->state,
34,136 ✔
1651
                                              VG_HANDLE_INFLIGHT,
1652
                                              VG_HANDLE_CB_DONE);
1653
  if (prev == VG_HANDLE_DRIVER_GONE) {
34,136 ✔
1654
    stTrace("stream vtable resolve cb: driver abandoned handle, cb destroys h=%p", h);
×
1655
    streamDestroyVgRpcHandle(&h);
×
1656
  }
1657

1658
  // Decrement shared pending; only the last completer posts the sem. The
1659
  // atomic dec is a release barrier -- the drain's matching tsem2_timewait
1660
  // acquire pairs with it so rsp/code stores above are visible after the wait.
1661
  int32_t remaining = atomic_sub_fetch_32(&pSync->pending, 1);
34,136 ✔
1662
  stTrace("stream vtable resolve cb done: code=0x%x remaining=%d", code, remaining);
34,136 ✔
1663
  if (remaining == 0) {
34,136 ✔
1664
    TAOS_UNUSED(tsem2_post(&pSync->sem));
23,404 ✔
1665
  }
1666

1667
  // Release this cb's lifetime ref on the sync. If this was the last ref
1668
  // (driver already walked away after timeout AND we are the last in-flight
1669
  // cb), this frees the sync. After this call pSync must not be touched.
1670
  TAOS_UNUSED(streamFanoutSyncRelease(pSync));
34,136 ✔
1671
  return code;
34,136 ✔
1672
}
1673

1674
// Phase A of fan-out: build a table-grouped request from the (batch, indexList)
1675
// slice, serialize it with SMsgHead, and asyncSendMsgToServer it.
1676
//
1677
// The caller owns `h` (already pre-allocated and pushed into the drain's
1678
// rpcHandles array) with h->vgId / h->indexList / h->ctx.pSync pre-filled.
1679
// This function only fills in h->req / h->totalCols and fires the RPC.
1680
//
1681
// On success returns 0; the callback is guaranteed to fire and decrement the
1682
// shared sync counter exactly once.
1683
// On ANY failure (serialize / send) returns the error code; h stays owned by
1684
// the caller and will be freed by the drain cleanup loop. The shared sync
1685
// counter is left untouched in the failure case (no callback will fire).
1686
static int32_t streamPrepareAndFireOneVgResolve(SVnode *pVnode, const SEpSet *pEpSet,
34,136 ✔
1687
                                                int64_t ver, SArray *batch, SArray *indexList,
1688
                                                SStreamVgRpcHandle *h) {
1689
  int32_t       code         = 0;
34,136 ✔
1690
  SHashObj     *tblGroupMap  = NULL;
34,136 ✔
1691
  void         *pReqBuf      = NULL;
34,136 ✔
1692
  SMsgSendInfo *pSendInfo    = NULL;
34,136 ✔
1693
  SArray       *groupPosList = NULL;  // SArray<SArray<int32_t>*>, temp per-group pos lists
34,136 ✔
1694

1695
  int32_t cnt = (int32_t)taosArrayGetSize(indexList);
34,136 ✔
1696
  stTrace("vgId:%d %s enter: targetVgId=%d ver=%" PRId64 " items=%d", TD_VID(pVnode), __func__,
34,136 ✔
1697
          h->vgId, ver, cnt);
1698

1699
  h->req.ver    = ver;
34,136 ✔
1700
  h->req.groups = taosArrayInit(4, sizeof(SVTableRefResolveGroupItem));
34,136 ✔
1701
  if (h->req.groups == NULL) { code = terrno; goto _capture; }
34,136 ✔
1702

1703
  // scatterOrder[i] records the dedupItems position of the i-th column in the
1704
  // flattened request (groups[0].cols[0], groups[0].cols[1], ...,
1705
  // groups[1].cols[0], ...).  The server returns responses in this same
1706
  // flattened order, so scatter must use scatterOrder — not indexList — to map
1707
  // rsp.items[i] back to its original slot in dedupRspItems.
1708
  h->scatterOrder = taosArrayInit(cnt, sizeof(int32_t));
34,136 ✔
1709
  if (h->scatterOrder == NULL) { code = terrno; goto _capture; }
34,136 ✔
1710

1711
  // Use a temp hash to map "dbName\0tableName" -> index in req.groups
1712
  tblGroupMap = taosHashInit(16, taosGetDefaultHashFunction(TSDB_DATA_TYPE_BINARY), true, HASH_NO_LOCK);
34,136 ✔
1713
  if (tblGroupMap == NULL) { code = terrno; goto _capture; }
34,136 ✔
1714

1715
  // Two-pass grouping: iterate indexList to build per-table groups (re-grouping
1716
  // interleaved columns like [t1.c1, t2.c1, t1.c2] into contiguous groups).
1717
  // Track per-group column positions in groupPosList so we can flatten them in
1718
  // group order into scatterOrder, which mirrors the server's response order.
1719
  groupPosList = taosArrayInit(4, sizeof(SArray *));
34,136 ✔
1720
  if (groupPosList == NULL) { code = terrno; goto _capture; }
34,136 ✔
1721

1722
  int32_t totalCols = 0;
34,136 ✔
1723
  for (int32_t i = 0; i < cnt; ++i) {
86,698 ✔
1724
    int32_t           pos = *(int32_t *)taosArrayGet(indexList, i);
52,629 ✔
1725
    SResolveWorkItem *w   = taosArrayGet(batch, pos);
52,629 ✔
1726

1727
    char    tblKey[TSDB_DB_NAME_LEN + 1 + TSDB_TABLE_NAME_LEN];
52,562 ✔
1728
    int32_t dLen = (int32_t)strlen(w->refDbName);
52,629 ✔
1729
    int32_t tLen = (int32_t)strlen(w->refTableName);
52,629 ✔
1730
    memcpy(tblKey, w->refDbName, dLen);
52,629 ✔
1731
    tblKey[dLen] = '\0';
52,629 ✔
1732
    memcpy(tblKey + dLen + 1, w->refTableName, tLen);
52,629 ✔
1733
    int32_t keyLen = dLen + 1 + tLen;
52,629 ✔
1734

1735
    int32_t *pGroupIdx = taosHashGet(tblGroupMap, tblKey, keyLen);
52,629 ✔
1736
    int32_t  groupIdx;
52,629 ✔
1737
    if (pGroupIdx == NULL) {
52,629 ✔
1738
      SVTableRefResolveGroupItem g = {0};
39,044 ✔
1739
      tstrncpy(g.dbName, w->refDbName, TSDB_DB_NAME_LEN);
39,044 ✔
1740
      tstrncpy(g.tableName, w->refTableName, TSDB_TABLE_NAME_LEN);
39,044 ✔
1741
      g.cols = taosArrayInit(4, sizeof(SVTableRefResolveColSpec));
39,044 ✔
1742
      if (g.cols == NULL) { code = terrno; goto _capture; }
39,044 ✔
1743
      if (taosArrayPush(h->req.groups, &g) == NULL) {
78,088 ✔
1744
        taosArrayDestroy(g.cols);
×
1745
        code = terrno;
×
1746
        goto _capture;
×
1747
      }
1748
      groupIdx = (int32_t)taosArrayGetSize(h->req.groups) - 1;
39,044 ✔
1749
      // Abort on put failure: a duplicate group would break the rsp scatter
1750
      // ordering assumption used by streamScatterOneVgResolve below.
1751
      if (taosHashPut(tblGroupMap, tblKey, keyLen, &groupIdx, sizeof(groupIdx)) != 0) {
39,044 ✔
1752
        code = terrno;
×
1753
        goto _capture;
×
1754
      }
1755
      // Create a parallel position list for this new group.
1756
      SArray *posList = taosArrayInit(4, sizeof(int32_t));
39,044 ✔
1757
      if (posList == NULL) { code = terrno; goto _capture; }
39,044 ✔
1758
      if (taosArrayPush(groupPosList, &posList) == NULL) {
39,044 ✔
1759
        taosArrayDestroy(posList);
×
1760
        code = terrno;
×
1761
        goto _capture;
×
1762
      }
1763
    } else {
1764
      groupIdx = *pGroupIdx;
13,585 ✔
1765
    }
1766

1767
    SVTableRefResolveGroupItem *gp = taosArrayGet(h->req.groups, groupIdx);
52,629 ✔
1768
    SVTableRefResolveColSpec    colSpec = {0};
52,629 ✔
1769
    tstrncpy(colSpec.colName, w->refColName, TSDB_COL_NAME_LEN);
52,629 ✔
1770
    colSpec.kind = w->kind;
52,629 ✔
1771
    if (taosArrayPush(gp->cols, &colSpec) == NULL) {
105,258 ✔
1772
      code = terrno;
×
1773
      goto _capture;
×
1774
    }
1775
    // Record the dedupItems position in the per-group list; this mirrors the
1776
    // column append order exactly, matching the server's response order.
1777
    SArray *curPosList = *(SArray **)taosArrayGet(groupPosList, groupIdx);
52,629 ✔
1778
    if (taosArrayPush(curPosList, &pos) == NULL) {
52,562 ✔
1779
      code = terrno;
×
1780
      goto _capture;
×
1781
    }
1782
    totalCols++;
52,562 ✔
1783
  }
1784

1785
  // Flatten per-group position lists into scatterOrder in group order.
1786
  // The server iterates groups[0], groups[1], ... and within each group
1787
  // appends cols[0], cols[1], ..., so this produces the exact same order.
1788
  int32_t nGroups = (int32_t)taosArrayGetSize(groupPosList);
34,069 ✔
1789
  for (int32_t g = 0; g < nGroups; ++g) {
73,180 ✔
1790
    SArray *posList = *(SArray **)taosArrayGet(groupPosList, g);
39,044 ✔
1791
    int32_t nPos    = (int32_t)taosArrayGetSize(posList);
39,044 ✔
1792
    for (int32_t j = 0; j < nPos; ++j) {
91,673 ✔
1793
      int32_t p = *(int32_t *)taosArrayGet(posList, j);
52,629 ✔
1794
      if (taosArrayPush(h->scatterOrder, &p) == NULL) {
105,258 ✔
1795
        code = terrno;
×
1796
        goto _capture;
×
1797
      }
1798
    }
1799
  }
1800
  // Free the temporary per-group position lists.
1801
  for (int32_t g = 0; g < nGroups; ++g) {
73,180 ✔
1802
    SArray *posList = *(SArray **)taosArrayGet(groupPosList, g);
39,044 ✔
1803
    taosArrayDestroy(posList);
39,044 ✔
1804
  }
1805
  taosArrayDestroy(groupPosList);
34,136 ✔
1806
  groupPosList = NULL;
34,136 ✔
1807

1808
  h->totalCols = totalCols;
34,136 ✔
1809

1810
  void *clientRpc = pVnode->msgCb.clientRpc;
34,136 ✔
1811
  if (clientRpc == NULL) { code = TSDB_CODE_INVALID_PARA; goto _capture; }
34,136 ✔
1812

1813
  int32_t reqLen = tSerializeSVTableRefResolveReq(NULL, 0, &h->req);
34,136 ✔
1814
  if (reqLen < 0) { code = terrno; goto _capture; }
34,136 ✔
1815
  // Prepend SMsgHead so dnode dispatcher (vmPutMsgToQueue) can route by vgId.
1816
  int32_t totalLen = reqLen + (int32_t)sizeof(SMsgHead);
34,136 ✔
1817
  pReqBuf = taosMemoryCalloc(1, totalLen);
34,136 ✔
1818
  if (pReqBuf == NULL) { code = terrno; goto _capture; }
34,136 ✔
1819
  if (tSerializeSVTableRefResolveReq((char *)pReqBuf + sizeof(SMsgHead), reqLen, &h->req) < 0) {
34,136 ✔
1820
    code = terrno;
×
1821
    goto _capture;
×
1822
  }
1823
  ((SMsgHead *)pReqBuf)->vgId    = htonl(h->vgId);
34,136 ✔
1824
  ((SMsgHead *)pReqBuf)->contLen = htonl(totalLen);
34,136 ✔
1825

1826
  pSendInfo = taosMemoryCalloc(1, sizeof(SMsgSendInfo));
34,136 ✔
1827
  if (pSendInfo == NULL) { code = terrno; goto _capture; }
34,136 ✔
1828

1829
  pSendInfo->param         = &h->ctx;
34,136 ✔
1830
  pSendInfo->msgInfo.pData = pReqBuf;
34,136 ✔
1831
  pSendInfo->msgInfo.len   = totalLen;
34,136 ✔
1832
  pSendInfo->msgType       = TDMT_VND_VTABLE_REF_RESOLVE;
34,136 ✔
1833
  pSendInfo->fp            = streamProcessVgResolveRsp;
34,136 ✔
1834
  pReqBuf = NULL;  // ownership transferred to pSendInfo
34,136 ✔
1835

1836
  // Reserve BOTH counters BEFORE asyncSend so a fast/synchronous callback
1837
  // cannot race past us. `pending` gates the sem post; `refs` gates sync
1838
  // lifetime. Both are rolled back if send is rejected.
1839
  TAOS_UNUSED(atomic_add_fetch_32(&h->ctx.pSync->pending, 1));
34,136 ✔
1840
  TAOS_UNUSED(atomic_add_fetch_32(&h->ctx.pSync->refs, 1));
34,136 ✔
1841

1842
  code = asyncSendMsgToServer(clientRpc, (SEpSet *)pEpSet, NULL, pSendInfo);
34,136 ✔
1843
  pSendInfo = NULL;  // ownership transferred (or freed by asyncSendMsgToServer on error)
34,136 ✔
1844
  stTrace("vgId:%d %s asyncSend done: targetVgId=%d code=0x%x reqLen=%d", TD_VID(pVnode), __func__,
34,136 ✔
1845
          h->vgId, code, totalLen);
1846
  if (code != 0) {
34,136 ✔
1847
    // Send rejected -- no callback will fire, release both reservations.
1848
    // The drain holds its own reservation on both counters, so neither can
1849
    // cross 0 here.
1850
    TAOS_UNUSED(atomic_sub_fetch_32(&h->ctx.pSync->pending, 1));
×
1851
    TAOS_UNUSED(atomic_sub_fetch_32(&h->ctx.pSync->refs, 1));
×
1852
    goto _capture;
×
1853
  }
1854

1855
  // Successfully queued: callback will fire and decrement the shared counter.
1856
  if (tblGroupMap != NULL) taosHashCleanup(tblGroupMap);
34,136 ✔
1857
  return 0;
34,136 ✔
1858

1859
_capture:
×
1860
  // Any failure before asyncSend accepted the request: no callback will fire.
1861
  // h is owned by the caller (already in rpcHandles); free only local buffers.
1862
  if (tblGroupMap != NULL) taosHashCleanup(tblGroupMap);
×
1863
  if (pReqBuf != NULL) taosMemoryFree(pReqBuf);
×
1864
  if (pSendInfo != NULL) taosMemoryFree(pSendInfo);
×
1865
  if (groupPosList != NULL) {
×
1866
    int32_t ngl = (int32_t)taosArrayGetSize(groupPosList);
×
1867
    for (int32_t g = 0; g < ngl; ++g) {
×
1868
      SArray *pl = *(SArray **)taosArrayGet(groupPosList, g);
×
1869
      taosArrayDestroy(pl);
×
1870
    }
1871
    taosArrayDestroy(groupPosList);
×
1872
  }
1873
  return code;
×
1874
}
1875

1876
// Phase B of fan-out: scatter a single completed RPC's rsp items into
1877
// outRspItems using the scatterOrder mapping built during phase A. Items'
1878
// tagData ownership is transferred from the rsp into outRspItems (rsp slots
1879
// NULLed) to match the existing scatter semantics used downstream by the
1880
// deep-copy step.
1881
//
1882
// Waiting for callback completion is NOT done here; it is centralized in
1883
// streamBatchFanoutDrain via the shared SStreamFanoutSync. By the time the
1884
// drain returns, every fired callback has completed (pending == 0) and the
1885
// rsp/code fields are safely published.
1886
//
1887
// Returns 0 on success; non-zero on rsp decode error or size mismatch. Caller
1888
// must still call streamDestroyVgRpcHandle to release ctx/req memory.
1889
static int32_t streamScatterOneVgResolve(SVnode *pVnode, SStreamVgRpcHandle *h,
34,136 ✔
1890
                                         SArray *outRspItems) {
1891
  if (h == NULL) return TSDB_CODE_INVALID_PARA;
34,136 ✔
1892

1893
  stTrace("vgId:%d %s scatter: targetVgId=%d ctxCode=0x%x rspItems=%d", TD_VID(pVnode), __func__,
34,136 ✔
1894
          h->vgId, h->ctx.code, h->ctx.rsp.items ? (int32_t)taosArrayGetSize(h->ctx.rsp.items) : 0);
1895

1896
  if (h->ctx.code != 0) return h->ctx.code;
34,136 ✔
1897

1898
  int32_t cnt = (int32_t)taosArrayGetSize(h->scatterOrder);
34,136 ✔
1899
  int32_t m   = (h->ctx.rsp.items != NULL) ? (int32_t)taosArrayGetSize(h->ctx.rsp.items) : 0;
34,136 ✔
1900
  // Both scatterOrder and rsp.items must match totalCols; cnt mismatch would
1901
  // cause an out-of-bounds read in the loop below.
1902
  if (cnt != m || m != h->totalCols) return TSDB_CODE_INVALID_MSG;
34,136 ✔
1903

1904
  // Move each rsp item back to its original dedupItems slot.
1905
  // scatterOrder[i] is the dedupItems position of the i-th column in the
1906
  // flattened request, which exactly matches the server response order
1907
  // (groups[0].cols[0], groups[0].cols[1], ..., groups[1].cols[0], ...).
1908
  // Using scatterOrder instead of indexList is correct even when columns of
1909
  // the same table are interleaved across multiple tables in the original
1910
  // vtable column list (e.g. [t1.c1, t2.c1, t1.c2]), where the grouping
1911
  // re-orders them but scatterOrder tracks the mapping precisely.
1912
  for (int32_t i = 0; i < cnt; ++i) {
86,765 ✔
1913
    int32_t                   pos = *(int32_t *)taosArrayGet(h->scatterOrder, i);
52,629 ✔
1914
    SVTableRefResolveRspItem *src = taosArrayGet(h->ctx.rsp.items, i);
52,629 ✔
1915
    SVTableRefResolveRspItem *dst = taosArrayGet(outRspItems, pos);
52,629 ✔
1916
    *dst = *src;
52,629 ✔
1917
    src->tagData = NULL;
52,629 ✔
1918
    src->tagLen  = 0;
52,629 ✔
1919
  }
1920
  return 0;
34,136 ✔
1921
}
1922

1923

1924
// Drive one resolution round for a heterogeneous batch: group work-items by the
1925
// target vgId of (refDbName, refTableName), issue one RPC per vg, and write
1926
// responses back to outRspItems in batch order.
1927
//
1928
// pCache (optional): caches db routing info across hops/uids to avoid hammering
1929
// mnode. NULL means no cache (every miss goes to mnode).
1930
//
1931
// outRspItems must be pre-sized with batch.size() default-zero entries; this
1932
// function fills them in place.
1933
//
1934
// Build the flat composite key "dbName\0tableName\0colName" used by BOTH
1935
// the Phase-1 dedup map and the tblRefCache. Tags and columns cannot share
1936
// a name within the same physical table, so kind does not enter the key.
1937
// Caller-provided buffer must hold at least
1938
// TSDB_DB_NAME_LEN + TSDB_TABLE_NAME_LEN + TSDB_COL_NAME_LEN + 4 bytes.
1939
// Example: db="mydb", tb="t1", col="voltage" => key="mydb\0t1\0voltage" (len=16)
1940
void streamBuildTblColKey(const char *db, const char *tb, const char *col,
300,503 ✔
1941
                                 char *out, int32_t *outLen) {
1942
  int32_t n = 0;
300,503 ✔
1943
  int32_t dbLen = (int32_t)strlen(db);
300,503 ✔
1944
  memcpy(out + n, db, dbLen); n += dbLen;
300,503 ✔
1945
  out[n++] = '\0';
300,503 ✔
1946
  int32_t tbLen = (int32_t)strlen(tb);
300,503 ✔
1947
  memcpy(out + n, tb, tbLen); n += tbLen;
300,503 ✔
1948
  out[n++] = '\0';
300,503 ✔
1949
  int32_t clLen = (int32_t)strlen(col);
300,503 ✔
1950
  memcpy(out + n, col, clLen); n += clLen;
300,503 ✔
1951
  *outLen = n;
300,503 ✔
1952
}
300,503 ✔
1953

1954
// Helper: look up tblRefCache for a resolved column. Returns pointer to cached
1955
// SVTableRefResolveRspItem or NULL if not cached.
1956
SVTableRefResolveRspItem *streamTblRefCacheLookup(SStreamVTableInfoCache *pCache,
121,076 ✔
1957
                                                          const char *dbName, const char *tableName,
1958
                                                          const char *colName, int8_t kind) {
1959
  (void)kind;  // tag and col cannot share a name within a table
22 ✔
1960
  if (pCache == NULL || pCache->tblRefCache == NULL) return NULL;
121,076 ✔
1961
  char    key[TSDB_DB_NAME_LEN + TSDB_TABLE_NAME_LEN + TSDB_COL_NAME_LEN + 4];
106,747 ✔
1962
  int32_t keyLen = 0;
106,747 ✔
1963
  streamBuildTblColKey(dbName, tableName, colName, key, &keyLen);
106,807 ✔
1964
  return (SVTableRefResolveRspItem *)taosHashGet(pCache->tblRefCache, key, keyLen);
106,807 ✔
1965
}
1966

1967
// Helper: insert a resolved column result into tblRefCache.
1968
void streamTblRefCacheInsert(SStreamVTableInfoCache *pCache,
98,677 ✔
1969
                                     const char *dbName, const char *tableName,
1970
                                     const char *colName, int8_t kind,
1971
                                     const SVTableRefResolveRspItem *pItem) {
1972
  (void)kind;  // tag and col cannot share a name within a table
14 ✔
1973
  if (pCache == NULL || pCache->tblRefCache == NULL) return;
98,677 ✔
1974
  char    key[TSDB_DB_NAME_LEN + TSDB_TABLE_NAME_LEN + TSDB_COL_NAME_LEN + 4];
84,552 ✔
1975
  int32_t keyLen = 0;
84,552 ✔
1976
  streamBuildTblColKey(dbName, tableName, colName, key, &keyLen);
84,552 ✔
1977
  // Store a deep copy including tagData so cache outlives the original rsp.
1978
  SVTableRefResolveRspItem copy = *pItem;
84,552 ✔
1979
  if (pItem->tagData != NULL && pItem->tagLen > 0) {
84,497 ✔
1980
    copy.tagData = taosMemoryMalloc(pItem->tagLen);
4,692 ✔
1981
    if (copy.tagData != NULL) {
4,692 ✔
1982
      memcpy(copy.tagData, pItem->tagData, pItem->tagLen);
4,692 ✔
1983
      copy.tagLen = pItem->tagLen;
4,692 ✔
1984
    } else {
1985
      copy.tagData = NULL;
×
1986
      copy.tagLen  = 0;
×
1987
    }
1988
  } else {
1989
    copy.tagData = NULL;
79,805 ✔
1990
    copy.tagLen  = 0;
79,805 ✔
1991
  }
1992
  if (taosHashPut(pCache->tblRefCache, key, keyLen, &copy, sizeof(copy)) != 0) {
84,497 ✔
1993
    // Put failed (likely OOM/rehash); release the freshly deep-copied tagData
1994
    // to avoid leaking it. Cache miss on next lookup will simply re-resolve.
1995
    stWarn("%s taosHashPut failed for col=%s, code=0x%x", __func__, colName, terrno);
57 ✔
1996
    taosMemoryFreeClear(copy.tagData);
2 ✔
1997
    copy.tagLen = 0;
2 ✔
1998
  }
1999
}
2000

2001
// (dedup map and tblRefCache share streamBuildTblColKey above; no extra alias.)
2002

2003

2004
// Phase 1 helper: deep-copy one rsp item from cache (or remote) into outRspItems[i].
2005
int32_t streamWriteRspItemDeepCopy(const SVTableRefResolveRspItem *src,
121,316 ✔
2006
                                          SVTableRefResolveRspItem *dst) {
2007
  *dst = *src;
121,316 ✔
2008
  if (src->tagData != NULL && src->tagLen > 0) {
121,316 ✔
2009
    dst->tagData = taosMemoryMalloc(src->tagLen);
13,134 ✔
2010
    if (dst->tagData == NULL) {
13,134 ✔
2011
      // Mirror the existing graceful-degrade behavior: keep dst->code,
2012
      // surface the missing payload by zeroing tagLen instead of failing
2013
      // the entire batch.
2014
      dst->tagLen = 0;
×
2015
      return terrno;
×
2016
    }
2017
    memcpy(dst->tagData, src->tagData, src->tagLen);
13,134 ✔
2018
  } else {
2019
    dst->tagData = NULL;
108,182 ✔
2020
    dst->tagLen  = 0;
108,182 ✔
2021
  }
2022
  return 0;
121,308 ✔
2023
}
2024

2025
// Phase 1 of streamCallResolveBatched: walk batch[]; for each item, either
2026
// resolve from tblRefCache (writing directly into outRspItems[i] and marking
2027
// origToDedupIdx[i]=-1) or deduplicate by (db,table,col) into dedupItems and
2028
// record origToDedupIdx[i] = dedup slot.
2029
int32_t streamBatchTryCacheAndDedup(SStreamVTableInfoCache *pCache, SArray *batch,
49,122 ✔
2030
                                           SArray *outRspItems, SHashObj *dedupMap,
2031
                                           SArray *dedupItems, int32_t *origToDedupIdx,
2032
                                           int32_t *pCacheHits) {
2033
  int32_t n = (int32_t)taosArrayGetSize(batch);
49,122 ✔
2034
  int32_t cacheHits = 0;
49,181 ✔
2035
  for (int32_t i = 0; i < n; ++i) {
170,438 ✔
2036
    SResolveWorkItem *w = taosArrayGet(batch, i);
121,316 ✔
2037

2038
    SVTableRefResolveRspItem *cached = streamTblRefCacheLookup(pCache, w->refDbName, w->refTableName,
121,324 ✔
2039
                                                               w->refColName, w->kind);
121,316 ✔
2040
    if (cached != NULL) {
121,316 ✔
2041
      // Cache hit: deep copy into the caller-owned outRspItems slot directly.
2042
      SVTableRefResolveRspItem *dst = taosArrayGet(outRspItems, i);
12,196 ✔
2043
      int32_t rc = streamWriteRspItemDeepCopy(cached, dst);
12,196 ✔
2044
      if (rc != 0) {
12,196 ✔
2045
        // OOM: treat as cache miss so this item falls through to RPC path.
2046
        stWarn("streamBatchTryCacheAndDedup deep copy failed (OOM), falling back to RPC: i=%d rc=0x%x", i, rc);
×
2047
        dst->tagLen = 0;
×
2048
        dst->tagData = NULL;
×
2049
        // Fall through to dedup path below.
2050
      } else {
2051
        origToDedupIdx[i] = -1;
12,196 ✔
2052
        cacheHits++;
12,196 ✔
2053
        continue;
12,196 ✔
2054
      }
2055
    }
2056

2057
    // Dedup key: tag and col share namespace within a table so kind is omitted.
2058
    char    dedupKey[TSDB_DB_NAME_LEN + TSDB_TABLE_NAME_LEN + TSDB_COL_NAME_LEN + 4];
109,120 ✔
2059
    int32_t dkLen = 0;
109,061 ✔
2060
    streamBuildTblColKey(w->refDbName, w->refTableName, w->refColName, dedupKey, &dkLen);
109,061 ✔
2061

2062
    int32_t *pExistIdx = (int32_t *)taosHashGet(dedupMap, dedupKey, dkLen);
109,120 ✔
2063
    if (pExistIdx != NULL) {
109,001 ✔
2064
      origToDedupIdx[i] = *pExistIdx;
10,398 ✔
2065
    } else {
2066
      int32_t newIdx = (int32_t)taosArrayGetSize(dedupItems);
98,603 ✔
2067
      if (taosArrayPush(dedupItems, w) == NULL) return terrno;
98,604 ✔
2068
      if (taosHashPut(dedupMap, dedupKey, dkLen, &newIdx, sizeof(newIdx)) != 0) return terrno;
98,604 ✔
2069
      origToDedupIdx[i] = newIdx;
98,603 ✔
2070
    }
2071
  }
2072
  *pCacheHits = cacheHits;
49,122 ✔
2073
  return 0;
49,122 ✔
2074
}
2075

2076
// Phase 2 of streamCallResolveBatched: route each dedup item to its target vg
2077
// (via the mnode cache or live lookup) and bucket the dedup-index into
2078
// vg2Idx[vgId]. Also records the SEpSet per vg in vg2Ep for the fan-out phase.
2079
static int32_t streamBatchRouteToVgs(SVnode *pVnode, SStreamVTableInfoCache *pCache,
45,628 ✔
2080
                                     SArray *dedupItems, SHashObj *vg2Idx, SHashObj *vg2Ep) {
2081
  int32_t acctId = 0;
45,628 ✔
2082
  if (sscanf(pVnode->config.dbname, "%d.", &acctId) != 1) return TSDB_CODE_INVALID_PARA;
45,628 ✔
2083

2084
  int32_t dedupN = (int32_t)taosArrayGetSize(dedupItems);
45,628 ✔
2085
  for (int32_t i = 0; i < dedupN; ++i) {
144,346 ✔
2086
    SResolveWorkItem *w = taosArrayGet(dedupItems, i);
98,718 ✔
2087

2088
    char dbFName[TSDB_DB_FNAME_LEN] = {0};
98,718 ✔
2089
    (void)tsnprintf(dbFName, sizeof(dbFName), "%d.%s", acctId, w->refDbName);
98,718 ✔
2090

2091
    SUseDbRsp *pRsp      = NULL;
98,718 ✔
2092
    bool       fromCache = false;
98,718 ✔
2093
    int32_t    rc = streamGetOrFetchDbVgInfo(pVnode, pCache, dbFName, &pRsp, &fromCache);
98,718 ✔
2094
    if (rc != 0) {
98,718 ✔
2095
      stError("vgId:%d %s uid=%" PRId64 " getDbVgInfo db=%s rc=0x%x -> propagate",
×
2096
              TD_VID(pVnode), __func__, w->originVtbUid, dbFName, rc);
2097
      return rc;
×
2098
    }
2099

2100
    int32_t vgId  = 0;
98,718 ✔
2101
    SEpSet  epSet = {0};
98,718 ✔
2102
    rc = streamRouteTableToVg(pRsp, dbFName, w->refTableName, &vgId, &epSet);
98,718 ✔
2103
    if (!fromCache && pCache == NULL) {
98,718 ✔
2104
      tFreeSUsedbRsp(pRsp);
14,176 ✔
2105
      taosMemoryFree(pRsp);
14,176 ✔
2106
    }
2107
    if (rc != 0) {
98,718 ✔
2108
      stError("vgId:%d %s uid=%" PRId64 " routeTableToVg db=%s tb=%s rc=0x%x -> propagate",
×
2109
              TD_VID(pVnode), __func__, w->originVtbUid, dbFName, w->refTableName, rc);
2110
      return rc;
×
2111
    }
2112

2113
    SArray **ppList = (SArray **)taosHashGet(vg2Idx, &vgId, sizeof(vgId));
98,718 ✔
2114
    SArray  *pList  = NULL;
98,718 ✔
2115
    if (ppList == NULL) {
98,718 ✔
2116
      pList = taosArrayInit(4, sizeof(int32_t));
57,100 ✔
2117
      if (pList == NULL) return terrno;
57,100 ✔
2118
      if (taosHashPut(vg2Idx, &vgId, sizeof(vgId), &pList, sizeof(pList)) != 0) {
57,100 ✔
2119
        taosArrayDestroy(pList);
×
2120
        return terrno;
×
2121
      }
2122
      if (taosHashPut(vg2Ep, &vgId, sizeof(vgId), &epSet, sizeof(epSet)) != 0) return terrno;
57,100 ✔
2123
    } else {
2124
      pList = *ppList;
41,618 ✔
2125
    }
2126
    if (taosArrayPush(pList, &i) == NULL) return terrno;
197,436 ✔
2127
  }
2128
  return 0;
45,628 ✔
2129
}
2130

2131
// Phase 3a sub-step: run the local vg's items synchronously via vnodeResolveOneHop.
2132
// Returns the first OOM (if any); other per-item errors are stashed in dst->code
2133
// without aborting the loop, matching the original behavior.
2134
static int32_t streamBatchExecuteLocalVg(SVnode *pVnode, SArray *dedupItems, SArray *pList,
22,964 ✔
2135
                                         SArray *dedupRspItems) {
2136
  int32_t cnt = (int32_t)taosArrayGetSize(pList);
22,964 ✔
2137
  for (int32_t j = 0; j < cnt; ++j) {
69,053 ✔
2138
    int32_t           pos = *(int32_t *)taosArrayGet(pList, j);
46,089 ✔
2139
    SResolveWorkItem *w   = taosArrayGet(dedupItems, pos);
46,089 ✔
2140
    SVTableRefResolveItem q = {0};
46,089 ✔
2141
    q.kind   = w->kind;
46,089 ✔
2142
    q.hasRef = true;
46,089 ✔
2143
    tstrncpy(q.refDbName,    w->refDbName,    TSDB_DB_NAME_LEN);
46,089 ✔
2144
    tstrncpy(q.refTableName, w->refTableName, TSDB_TABLE_NAME_LEN);
46,089 ✔
2145
    tstrncpy(q.refColName,   w->refColName,   TSDB_COL_NAME_LEN);
46,089 ✔
2146
    SVTableRefResolveRspItem *dst = taosArrayGet(dedupRspItems, pos);
46,089 ✔
2147
    int32_t one = vnodeResolveOneHop(pVnode, &q, dst);
46,089 ✔
2148
    if (one != 0) {
46,089 ✔
2149
      dst->code = one;
×
2150
      if (one == TSDB_CODE_OUT_OF_MEMORY) return one;
×
2151
    }
2152
  }
2153
  return 0;
22,964 ✔
2154
}
2155

2156
// Phase 3 of streamCallResolveBatched: fan-out one RPC per remote vg + drain.
2157
// See the header comment of streamCallResolveBatched for the concurrency model.
2158
//
2159
// Signalling: a single heap-allocated SStreamFanoutSync (binary sem + atomic
2160
// pending + atomic refs) is shared across every fired RPC handle.
2161
//   - pending gates the sem post: the cb that drives pending->0 posts the sem.
2162
//   - refs gates sync lifetime: driver holds 1, each successful fire holds 1,
2163
//     whichever side drives refs->0 destroys the sync.
2164
//
2165
// Wait uses tsem2_timewait. On timeout the driver atomically hands off each
2166
// still-INFLIGHT handle to its (future) cb via a per-handle state CAS, so
2167
// late cbs free both their handle and (when last) the shared sync without
2168
// touching driver-frame storage. The driver then releases its sync ref and
2169
// returns TSDB_CODE_TIMEOUT_ERROR.
2170
static int32_t streamBatchFanoutDrain(SVnode *pVnode, int64_t ver, SArray *dedupItems,
45,628 ✔
2171
                                      SHashObj *vg2Idx, SHashObj *vg2Ep, SArray *dedupRspItems) {
2172
  int32_t            code        = 0;
45,628 ✔
2173
  void              *pIter       = NULL;
45,628 ✔
2174
  SStreamFanoutSync *pSync       = NULL;
45,628 ✔
2175
  SArray            *rpcHandles  = NULL;
45,628 ✔
2176

2177
  // Pre-reserve capacity equal to the number of remote vgs so taosArrayPush
2178
  // does not need to grow / allocate during fan-out.
2179
  int32_t nRemoteVgs = taosHashGetSize(vg2Idx);
45,628 ✔
2180
  if (nRemoteVgs < 1) nRemoteVgs = 1;
45,628 ✔
2181
  rpcHandles = taosArrayInit(nRemoteVgs, sizeof(SStreamVgRpcHandle *));
45,628 ✔
2182
  if (rpcHandles == NULL) {
45,628 ✔
2183
    code = terrno;
×
2184
    stError("vgId:%d %s init rpcHandles failed: code=0x%x", TD_VID(pVnode), __func__, code);
×
2185
    goto _exit;
×
2186
  }
2187

2188
  pSync = streamFanoutSyncCreate();
45,628 ✔
2189
  if (pSync == NULL) {
45,628 ✔
2190
    code = terrno;
×
2191
    stError("vgId:%d %s create sync failed: code=0x%x", TD_VID(pVnode), __func__, code);
×
2192
    goto _exit;
×
2193
  }
2194
  // Reserve a "fan-out in progress" slot so `pending` cannot transiently
2195
  // hit 0 between firing handles; released right before tsem2_timewait.
2196
  atomic_store_32(&pSync->pending, 1);
45,628 ✔
2197
  // Driver holds one lifetime ref; each successful fire adds one. Whichever
2198
  // party drives refs->0 calls streamFanoutSyncDestroy.
2199
  atomic_store_32(&pSync->refs, 1);
45,628 ✔
2200

2201
  // Phase 3a: process local-vg in-process; fire async RPC for each remote vg.
2202
  // On any error, stop iterating immediately and jump to cleanup; already
2203
  // fired RPCs are still drained in _exit to preserve callback ownership.
2204
  pIter = taosHashIterate(vg2Idx, NULL);
45,628 ✔
2205
  while (pIter != NULL) {
102,728 ✔
2206
    SArray  *pList  = *(SArray **)pIter;
57,100 ✔
2207
    size_t   keyLen = 0;
57,100 ✔
2208
    int32_t *pVgKey = (int32_t *)taosHashGetKey(pIter, &keyLen);
57,100 ✔
2209
    int32_t  vgId   = *pVgKey;
57,100 ✔
2210
    SEpSet  *pEpSet = (SEpSet *)taosHashGet(vg2Ep, &vgId, sizeof(vgId));
57,100 ✔
2211

2212
    if (vgId == TD_VID(pVnode)) {
57,100 ✔
2213
      int32_t rc = streamBatchExecuteLocalVg(pVnode, dedupItems, pList, dedupRspItems);
22,964 ✔
2214
      stTrace("vgId:%d %s local-vg done: targetVgId=%d items=%d rc=0x%x", TD_VID(pVnode),
22,964 ✔
2215
              __func__, vgId, (int32_t)taosArrayGetSize(pList), rc);
2216
      if (rc != 0) {
22,964 ✔
2217
        code = rc;
×
2218
        stError("vgId:%d %s local-vg failed: targetVgId=%d code=0x%x",
×
2219
                TD_VID(pVnode), __func__, vgId, code);
2220
        goto _exit;
×
2221
      }
2222
    } else {
2223
      // Pre-allocate and push the handle BEFORE firing. This way the only
2224
      // race-free failure modes are alloc/push (no RPC in flight) and
2225
      // prepare-fire (handle owned by array; cleanup loop frees it).
2226
      SStreamVgRpcHandle *h = taosMemoryCalloc(1, sizeof(SStreamVgRpcHandle));
34,136 ✔
2227
      if (h == NULL) {
34,136 ✔
2228
        code = terrno;
×
2229
        stError("vgId:%d %s alloc handle failed: targetVgId=%d code=0x%x",
×
2230
                TD_VID(pVnode), __func__, vgId, code);
2231
        goto _exit;
×
2232
      }
2233
      h->vgId      = vgId;
34,136 ✔
2234
      h->state     = VG_HANDLE_INFLIGHT;  // explicit; calloc already zeroed it
34,136 ✔
2235
      h->indexList = pList;
34,136 ✔
2236
      h->ctx.pSync = pSync;
34,136 ✔
2237

2238
      if (taosArrayPush(rpcHandles, &h) == NULL) {
34,136 ✔
2239
        // Capacity was pre-reserved above so this should not happen.
2240
        code = terrno;
×
2241
        stError("vgId:%d %s push handle failed: targetVgId=%d code=0x%x",
×
2242
                TD_VID(pVnode), __func__, vgId, code);
2243
        streamDestroyVgRpcHandle(&h);
×
2244
        goto _exit;
×
2245
      }
2246

2247
      int32_t rc = streamPrepareAndFireOneVgResolve(pVnode, pEpSet, ver,
34,136 ✔
2248
                                                    dedupItems, pList, h);
2249
      if (rc != 0) {
34,136 ✔
2250
        // h is already in rpcHandles; the cleanup loop in _exit frees it.
2251
        // No callback fired, so sync counters are untouched.
2252
        code = rc;
×
2253
        stError("vgId:%d %s prepare/fire failed: targetVgId=%d code=0x%x",
×
2254
                TD_VID(pVnode), __func__, vgId, code);
2255
        goto _exit;
×
2256
      }
2257
    }
2258
    pIter = taosHashIterate(vg2Idx, pIter);
57,100 ✔
2259
  }
2260

2261
_exit:
45,628 ✔
2262
  if (pIter != NULL) {
45,628 ✔
2263
    taosHashCancelIterate(vg2Idx, pIter);
×
2264
    pIter = NULL;
×
2265
  }
2266

2267
  // Phase 3b: release the driver's pending reservation, then either skip
2268
  // the wait (nothing in flight), wait normally, or time out and hand off.
2269
  if (pSync != NULL) {
45,628 ✔
2270
    int32_t remaining = atomic_sub_fetch_32(&pSync->pending, 1);
45,628 ✔
2271
    stTrace("vgId:%d %s drain release: remaining=%d", TD_VID(pVnode), __func__, remaining);
45,628 ✔
2272
    if (remaining != 0) {
45,628 ✔
2273
      int32_t waitRc = tsem2_timewait(&pSync->sem, STREAM_VTB_RPC_TIMEOUT_MS);
23,404 ✔
2274
      if (waitRc != 0) {
23,404 ✔
2275
        // Real timeout (a post racing exactly with the deadline would have
2276
        // been consumed by timewait itself). Atomically abandon every still
2277
        // in-flight handle. Each successful CAS transfers ownership of that
2278
        // handle to its (future) cb.
2279
        stWarn("vgId:%d %s fan-out timed out after %dms; abandoning in-flight handles",
×
2280
               TD_VID(pVnode), __func__, STREAM_VTB_RPC_TIMEOUT_MS);
2281
        if (rpcHandles != NULL) {
×
2282
          int32_t nHandles = (int32_t)taosArrayGetSize(rpcHandles);
×
2283
          for (int32_t i = 0; i < nHandles; ++i) {
×
2284
            SStreamVgRpcHandle **pSlot = (SStreamVgRpcHandle **)taosArrayGet(rpcHandles, i);
×
2285
            SStreamVgRpcHandle  *h     = *pSlot;
×
2286
            if (h == NULL) continue;
×
2287
            int8_t prev = atomic_val_compare_exchange_8(&h->state,
×
2288
                                                        VG_HANDLE_INFLIGHT,
2289
                                                        VG_HANDLE_DRIVER_GONE);
2290
            if (prev == VG_HANDLE_INFLIGHT) {
×
2291
              // Handoff succeeded: cb will free this handle when it eventually
2292
              // fires. Null the slot so the cleanup loop below skips it.
2293
              *pSlot = NULL;
×
2294
            }
2295
            // else: cb already completed (state == VG_HANDLE_CB_DONE); the
2296
            // driver still owns h and the cleanup loop will destroy it.
2297
          }
2298
        }
2299
        code = TSDB_CODE_TIMEOUT_ERROR;
×
2300
      }
2301
    }
2302
  }
2303

2304
  // Phase 3c: on success scatter rsp from each handle; on error skip the
2305
  // scatter (results are not consumable) and just destroy handles owned
2306
  // by the driver. NULL slots (abandoned on timeout) are owned by cbs.
2307
  if (rpcHandles != NULL) {
45,628 ✔
2308
    int32_t nHandles = (int32_t)taosArrayGetSize(rpcHandles);
45,628 ✔
2309
    for (int32_t i = 0; i < nHandles; ++i) {
79,709 ✔
2310
      SStreamVgRpcHandle **pSlot = (SStreamVgRpcHandle **)taosArrayGet(rpcHandles, i);
34,136 ✔
2311
      SStreamVgRpcHandle  *h     = *pSlot;
34,136 ✔
2312
      if (h == NULL) continue;
34,136 ✔
2313
      if (code == 0) {
34,136 ✔
2314
        int32_t rc = streamScatterOneVgResolve(pVnode, h, dedupRspItems);
34,136 ✔
2315
        stTrace("vgId:%d %s remote-vg done: targetVgId=%d items=%d rc=0x%x", TD_VID(pVnode),
34,136 ✔
2316
                __func__, h->vgId, (int32_t)taosArrayGetSize(h->indexList), rc);
2317
        if (rc != 0) {
34,136 ✔
2318
          code = rc;
×
2319
          stError("vgId:%d %s scatter failed: targetVgId=%d code=0x%x",
×
2320
                  TD_VID(pVnode), __func__, h->vgId, code);
2321
        }
2322
      }
2323
      streamDestroyVgRpcHandle(&h);
34,136 ✔
2324
      *pSlot = NULL;
34,081 ✔
2325
    }
2326
    taosArrayDestroy(rpcHandles);
45,573 ✔
2327
  }
2328

2329
  // Release the driver's lifetime ref on the sync. If any cbs are still
2330
  // pending (timeout path), the last one frees it; otherwise this call
2331
  // does. After this point pSync must not be touched.
2332
  TAOS_UNUSED(streamFanoutSyncRelease(pSync));
45,628 ✔
2333
  return code;
45,573 ✔
2334
}
2335

2336
// Phase 4 of streamCallResolveBatched: publish per-(table,col) results into
2337
// the tblRefCache (so future hops in this resolve cycle see them) and scatter
2338
// dedup results back to the caller's outRspItems[] positions, deep-copying
2339
// tagData since one dedup slot may fan out to multiple original positions.
2340
static int32_t streamBatchScatterAndPublish(SStreamVTableInfoCache *pCache, SArray *dedupItems,
45,573 ✔
2341
                                            SArray *dedupRspItems, int32_t *origToDedupIdx,
2342
                                            int32_t n, SArray *outRspItems) {
2343
  int32_t dedupN = (int32_t)taosArrayGetSize(dedupItems);
45,573 ✔
2344
  for (int32_t i = 0; i < dedupN; ++i) {
144,291 ✔
2345
    SResolveWorkItem         *w   = taosArrayGet(dedupItems, i);
98,663 ✔
2346
    SVTableRefResolveRspItem *rsp = taosArrayGet(dedupRspItems, i);
98,663 ✔
2347
    streamTblRefCacheInsert(pCache, w->refDbName, w->refTableName, w->refColName, w->kind, rsp);
98,718 ✔
2348
  }
2349

2350
  for (int32_t i = 0; i < n; ++i) {
154,742 ✔
2351
    int32_t dedupIdx = origToDedupIdx[i];
109,114 ✔
2352
    if (dedupIdx < 0) continue;  // already filled from cache in Phase 1
109,114 ✔
2353
    SVTableRefResolveRspItem *src = taosArrayGet(dedupRspItems, dedupIdx);
109,114 ✔
2354
    SVTableRefResolveRspItem *dst = taosArrayGet(outRspItems, i);
109,114 ✔
2355
    int32_t rc = streamWriteRspItemDeepCopy(src, dst);
109,114 ✔
2356
    if (rc != 0) {
109,114 ✔
2357
      // OOM during scatter: dst has incomplete tagData; propagate error so
2358
      // the caller can fail the batch rather than return corrupted results.
2359
      return rc;
×
2360
    }
2361
  }
2362
  return 0;
45,628 ✔
2363
}
2364

2365
//
2366
// streamCallResolveBatched: drive one hop of resolution with table-level dedup.
2367
//
2368
// Optimization (Issue 4): instead of sending per-(table,column) items blindly,
2369
// we (a) check the local tblRefCache first, (b) deduplicate by (db,table,col)
2370
// so the same physical column is only resolved once per RPC round, and (c) cache
2371
// the results for use in subsequent hops.
2372
//
2373
// Concurrency (review finding #4): remote per-vg RPCs are FAN-OUT — phase 3a
2374
// fires all asyncSends in a tight loop, phase 3b drains every fired handle.
2375
// Total wall time is bounded by max(per-vg RTT) instead of sum(per-vg RTT).
2376
// The local vg (if present) is still processed synchronously in phase 3a since
2377
// it bypasses RPC entirely.
2378
//
2379
// H2 v0.5 strict: any per-vg routing/RPC failure is propagated upward as the
2380
// return value. Per-item business errors are still reported through
2381
// outRspItems[i].code so the caller can include the originating uid/cid in
2382
// its log; the caller (streamResolveVTableRefChain) decides how to react.
2383
//
2384
// Returns 0 on success; non-zero on OOM, routing, or RPC failure.
2385
//
2386
// This driver is intentionally thin — each phase lives in its own helper
2387
// (streamBatchTryCacheAndDedup / streamBatchRouteToVgs / streamBatchFanoutDrain
2388
// / streamBatchScatterAndPublish) so the per-phase invariants are easier to
2389
// read and modify in isolation.
2390
static int32_t streamCallResolveBatched(SVnode *pVnode, SStreamVTableInfoCache *pCache,
49,179 ✔
2391
                                        int64_t ver, SArray *batch, SArray *outRspItems) {
2392
  int32_t   code     = 0;
49,179 ✔
2393
  SHashObj *vg2Idx   = NULL;  // key: int32_t vgId, value: SArray<int32_t>* (positions in dedupItems)
49,179 ✔
2394
  SHashObj *vg2Ep    = NULL;  // key: int32_t vgId, value: SEpSet
49,179 ✔
2395
  SHashObj *dedupMap = NULL;  // key: "db\0table\0col", value: int32_t (position in dedupItems)
49,179 ✔
2396
  SArray   *dedupItems    = NULL;  // SArray<SResolveWorkItem> unique items to send
49,179 ✔
2397
  SArray   *dedupRspItems = NULL;  // SArray<SVTableRefResolveRspItem> responses for dedup items
49,179 ✔
2398
  int32_t  *origToDedupIdx = NULL;  // batch index -> dedupItems index, or -1 if served from cache
49,179 ✔
2399

2400
  int32_t n = (int32_t)taosArrayGetSize(batch);
49,179 ✔
2401
  stDebug("vgId:%d %s enter: ver=%" PRId64 " batch=%d", TD_VID(pVnode), __func__, ver, n);
49,179 ✔
2402
  // Pre-size outRspItems with n zero entries so positional writes are safe.
2403
  for (int32_t i = (int32_t)taosArrayGetSize(outRspItems); i < n; ++i) {
170,487 ✔
2404
    SVTableRefResolveRspItem zero = {0};
121,308 ✔
2405
    if (taosArrayPush(outRspItems, &zero) == NULL) { code = terrno; goto _end; }
121,308 ✔
2406
  }
2407

2408
  // ---- Phase 1: cache lookup + dedup ----
2409
  dedupMap   = taosHashInit(n, taosGetDefaultHashFunction(TSDB_DATA_TYPE_BINARY), false, HASH_NO_LOCK);
49,179 ✔
2410
  dedupItems = taosArrayInit(n, sizeof(SResolveWorkItem));
49,119 ✔
2411
  if (dedupMap == NULL || dedupItems == NULL) { code = terrno; goto _end; }
49,179 ✔
2412
  origToDedupIdx = taosMemoryCalloc(n, sizeof(int32_t));
49,179 ✔
2413
  if (origToDedupIdx == NULL) { code = terrno; goto _end; }
49,179 ✔
2414

2415
  int32_t cacheHits = 0;
49,179 ✔
2416
  code = streamBatchTryCacheAndDedup(pCache, batch, outRspItems, dedupMap, dedupItems,
49,179 ✔
2417
                                     origToDedupIdx, &cacheHits);
2418
  if (code != 0) goto _end;
49,179 ✔
2419

2420
  int32_t dedupN = (int32_t)taosArrayGetSize(dedupItems);
49,179 ✔
2421
  stDebug("vgId:%d %s dedup: batch=%d cacheHits=%d dedupItems=%d",
49,179 ✔
2422
          TD_VID(pVnode), __func__, n, cacheHits, dedupN);
2423
  if (dedupN == 0) goto _end;  // all items served from cache
49,179 ✔
2424

2425
  // ---- Phase 2: route dedup items to vg groups ----
2426
  vg2Idx = taosHashInit(8, taosGetDefaultHashFunction(TSDB_DATA_TYPE_INT), false, HASH_NO_LOCK);
45,628 ✔
2427
  vg2Ep  = taosHashInit(8, taosGetDefaultHashFunction(TSDB_DATA_TYPE_INT), false, HASH_NO_LOCK);
45,628 ✔
2428
  if (vg2Idx == NULL || vg2Ep == NULL) { code = terrno; goto _end; }
45,628 ✔
2429
  code = streamBatchRouteToVgs(pVnode, pCache, dedupItems, vg2Idx, vg2Ep);
45,628 ✔
2430
  if (code != 0) goto _end;
45,628 ✔
2431

2432
  // ---- Phase 3: fan-out + drain ----
2433
  dedupRspItems = taosArrayInit(dedupN, sizeof(SVTableRefResolveRspItem));
45,628 ✔
2434
  if (dedupRspItems == NULL) { code = terrno; goto _end; }
45,628 ✔
2435
  for (int32_t i = 0; i < dedupN; ++i) {
144,346 ✔
2436
    SVTableRefResolveRspItem zero = {0};
98,718 ✔
2437
    if (taosArrayPush(dedupRspItems, &zero) == NULL) { code = terrno; goto _end; }
98,718 ✔
2438
  }
2439
  code = streamBatchFanoutDrain(pVnode, ver, dedupItems, vg2Idx, vg2Ep, dedupRspItems);
45,628 ✔
2440
  if (code != 0) {
45,573 ✔
2441
    stError("vgId:%d %s fan-out failed: rc=0x%x", TD_VID(pVnode), __func__, code);
×
2442
    goto _end;
×
2443
  }
2444

2445
  // ---- Phase 4: publish to cache + scatter to outRspItems ----
2446
  code = streamBatchScatterAndPublish(pCache, dedupItems, dedupRspItems, origToDedupIdx, n, outRspItems);
45,573 ✔
2447
  if (code != 0) {
45,628 ✔
2448
    stError("vgId:%d %s scatter failed (OOM): rc=0x%x", TD_VID(pVnode), __func__, code);
×
2449
  }
2450

2451
_end:
49,179 ✔
2452
  stDebug("vgId:%d %s exit: code=0x%x", TD_VID(pVnode), __func__, code);
49,179 ✔
2453
  taosMemoryFreeClear(origToDedupIdx);
49,179 ✔
2454
  if (vg2Idx != NULL) {
49,179 ✔
2455
    void *p = taosHashIterate(vg2Idx, NULL);
45,628 ✔
2456
    while (p != NULL) {
102,728 ✔
2457
      taosArrayDestroy(*(SArray **)p);
57,100 ✔
2458
      p = taosHashIterate(vg2Idx, p);
57,040 ✔
2459
    }
2460
    taosHashCleanup(vg2Idx);
45,628 ✔
2461
  }
2462
  if (vg2Ep != NULL) taosHashCleanup(vg2Ep);
49,119 ✔
2463
  if (dedupMap != NULL) taosHashCleanup(dedupMap);
49,059 ✔
2464
  taosArrayDestroy(dedupItems);
49,059 ✔
2465
  if (dedupRspItems != NULL) {
49,179 ✔
2466
    // Free any remaining tagData in dedup responses that were not transferred
2467
    int32_t sz = (int32_t)taosArrayGetSize(dedupRspItems);
45,628 ✔
2468
    for (int32_t i = 0; i < sz; ++i) {
144,286 ✔
2469
      SVTableRefResolveRspItem *r = taosArrayGet(dedupRspItems, i);
98,658 ✔
2470
      taosMemoryFreeClear(r->tagData);
98,658 ✔
2471
    }
2472
    taosArrayDestroy(dedupRspItems);
45,628 ✔
2473
  }
2474
  return code;
49,119 ✔
2475
}
2476

2477
// Function A: drive multi-hop chain resolution for a batch of vtable uids on the
2478
// triggering vnode. Cross-vgId version: groups each batch by target vgId, then
2479
// dispatches one TDMT_VND_VTABLE_REF_RESOLVE RPC per group via streamCallResolveBatched.
2480
//
2481
// H2 v0.5 strict error policy:
2482
//   - top-level uid not in local meta (or not a vtable type) -> warn + skip
2483
//     that uid; function returns 0 and uid simply has no entry in *ppUid2Result.
2484
//   - any other error (mid-chain table/col/tag missing, RPC failure, OOM,
2485
//     hop > MAX_HOPS, ref-triple inconsistency) -> A returns the underlying
2486
//     errCode; caller (reader -> trigger -> mnode) propagates and fail-fasts.
2487
// pCache (optional): caches db routing info (SUseDbRsp) across calls.
2488
// pReaderInfo (optional): when vtbUids is NULL/empty, all live uids are pulled from
2489
//                          qStreamGetTableArrayList(pReaderInfo). If both are NULL/empty
2490
//                          this function returns INVALID_PARA.
2491
// Output: *ppUid2Result is a fresh SSHashObj<uid -> SVTableResolveResult*>;
2492
// caller owns it and must use streamVTableResolveResultDestroy + tSimpleHashCleanup.
2493
// Full-uid branch helper: pull live (non-deleted) uids from the reader's table
2494
// list into a newly-allocated SArray<int64_t>. The caller owns both *ppFullUids
2495
// and *ppTableListArray and must free them. On failure both outputs are NULL.
2496
static int32_t streamCollectActiveVtableUids(SStreamTriggerReaderInfo *pReaderInfo,
4,269 ✔
2497
                                             SArray **ppTableListArray, SArray **ppFullUids) {
2498
  *ppTableListArray = NULL;
4,269 ✔
2499
  *ppFullUids       = NULL;
4,269 ✔
2500
  if (pReaderInfo == NULL) return TSDB_CODE_INVALID_PARA;
4,269 ✔
2501

2502
  SArray *pTableListArray = qStreamGetTableArrayList(pReaderInfo);
4,269 ✔
2503
  if (pTableListArray == NULL) return terrno;
4,269 ✔
2504

2505
  int32_t nAll     = (int32_t)taosArrayGetSize(pTableListArray);
4,269 ✔
2506
  SArray *fullUids = taosArrayInit(nAll, sizeof(int64_t));
4,269 ✔
2507
  if (fullUids == NULL) {
4,269 ✔
2508
    taosArrayDestroyP(pTableListArray, taosMemFree);
×
2509
    return terrno;
×
2510
  }
2511
  for (int32_t i = 0; i < nAll; ++i) {
10,862 ✔
2512
    SStreamTableKeyInfo *pKey = taosArrayGetP(pTableListArray, i);
6,593 ✔
2513
    if (pKey == NULL || pKey->markedDeleted) continue;
6,593 ✔
2514
    if (taosArrayPush(fullUids, &pKey->uid) == NULL) {
13,186 ✔
2515
      taosArrayDestroy(fullUids);
×
2516
      taosArrayDestroyP(pTableListArray, taosMemFree);
×
2517
      return terrno;
×
2518
    }
2519
  }
2520
  taosArrayRemoveDuplicate(fullUids, compareInt64Val, NULL);
4,269 ✔
2521
  *ppTableListArray = pTableListArray;
4,269 ✔
2522
  *ppFullUids       = fullUids;
4,269 ✔
2523
  return 0;
4,269 ✔
2524
}
2525

2526
// Consume one hop's rspItems[]: for each (workItem, rspItem) pair, either
2527
// propagate the per-item error, materialize the terminated result into
2528
// uid2Result (colMap or tagMap), or push the follow-up work item into
2529
// nextWorkList. This helper takes ownership of every rspItem.tagData (frees
2530
// it on every path) so the caller only has to manage the array itself.
2531
static int32_t streamConsumeOneHopResults(SVnode *pVnode, int32_t hop, SArray *workList,
49,179 ✔
2532
                                          SArray *rspItems, SSHashObj *uid2Result,
2533
                                          SArray *nextWorkList) {
2534
  int32_t bn = (int32_t)taosArrayGetSize(workList);
49,179 ✔
2535
  for (int32_t i = 0; i < bn; ++i) {
170,353 ✔
2536
    SResolveWorkItem         *w = taosArrayGet(workList, i);
121,308 ✔
2537
    SVTableRefResolveRspItem *r = taosArrayGet(rspItems, i);
121,308 ✔
2538

2539
    if (r->code != 0) {
121,308 ✔
2540
      // H2 v0.5: any per-item business error (mid-chain ref-table missing,
2541
      // ref-col missing, tag changed, etc.) is propagated upward.
2542
      stError("vgId:%d %s hop=%d uid=%" PRId64 " kind=%d cid=%d rspCode=0x%x -> propagate",
134 ✔
2543
              TD_VID(pVnode), __func__, hop, w->originVtbUid, w->kind, w->originCid, r->code);
2544
      int32_t rspCode = r->code;
134 ✔
2545
      taosMemoryFreeClear(r->tagData);
134 ✔
2546
      return rspCode;
134 ✔
2547
    }
2548

2549
    if (r->terminated) {
121,174 ✔
2550
      SVTableResolveResult *pRes = streamGetOrCreateUidResult(uid2Result, w->originVtbUid);
94,039 ✔
2551
      if (pRes == NULL) { taosMemoryFreeClear(r->tagData); return terrno; }
94,039 ✔
2552

2553
      if (w->kind == STREAM_VREF_KIND_COL) {
94,039 ✔
2554
        SColResolveItem *item = taosMemoryCalloc(1, sizeof(*item));
80,907 ✔
2555
        if (item == NULL) { taosMemoryFreeClear(r->tagData); return terrno; }
80,852 ✔
2556
        item->hasRef = r->nextRef.hasRef;
80,852 ✔
2557
        if (item->hasRef) {
80,852 ✔
2558
          tstrncpy(item->refDbName,    r->nextRef.refDbName,    TSDB_DB_NAME_LEN);
80,907 ✔
2559
          tstrncpy(item->refTableName, r->nextRef.refTableName, TSDB_TABLE_NAME_LEN);
80,907 ✔
2560
          tstrncpy(item->refColName,   r->nextRef.refColName,   TSDB_COL_NAME_LEN);
80,907 ✔
2561
        }
2562
        // Snapshot old pointer before put; free it only after successful put so
2563
        // the hash never holds a dangling pointer (avoids double-free on cleanup).
2564
        SColResolveItem **ppOld = (SColResolveItem **)tSimpleHashGet(pRes->colMap, &w->originCid, sizeof(w->originCid));
80,852 ✔
2565
        SColResolveItem  *oldItem = (ppOld && *ppOld) ? *ppOld : NULL;
80,907 ✔
2566
        if (tSimpleHashPut(pRes->colMap, &w->originCid, sizeof(w->originCid), &item, sizeof(item)) != 0) {
80,907 ✔
2567
          taosMemoryFree(item);
×
2568
          taosMemoryFreeClear(r->tagData);
×
2569
          return terrno;
×
2570
        }
2571
        if (oldItem) { taosMemoryFree(oldItem); }
80,907 ✔
2572
        stDebug("vgId:%d %s hop=%d uid=%" PRId64 " COL cid=%d TERMINATED hasRef=%d ref=%s.%s.%s -> colMap",
80,907 ✔
2573
                TD_VID(pVnode), __func__, hop, w->originVtbUid, w->originCid,
2574
                item->hasRef, item->refDbName, item->refTableName, item->refColName);
2575
        taosMemoryFreeClear(r->tagData);
80,907 ✔
2576
      } else {
2577
        STagValue *tv = taosMemoryCalloc(1, sizeof(*tv));
13,132 ✔
2578
        if (tv == NULL) { taosMemoryFreeClear(r->tagData); return terrno; }
13,132 ✔
2579
        tv->type  = r->tagType;
13,132 ✔
2580
        tv->nLen  = r->tagLen;
13,132 ✔
2581
        tv->pData = r->tagData;
13,132 ✔
2582
        r->tagData = NULL;  // ownership transferred to STagValue
13,132 ✔
2583
        // Snapshot old pointer before put; free it only after successful put so
2584
        // the hash never holds a dangling pointer (avoids double-free on cleanup).
2585
        STagValue **ppOldTag = (STagValue **)tSimpleHashGet(pRes->tagMap, &w->originCid, sizeof(w->originCid));
13,132 ✔
2586
        STagValue  *oldTag   = (ppOldTag && *ppOldTag) ? *ppOldTag : NULL;
13,132 ✔
2587
        if (tSimpleHashPut(pRes->tagMap, &w->originCid, sizeof(w->originCid), &tv, sizeof(tv)) != 0) {
13,132 ✔
2588
          taosMemoryFreeClear(tv->pData);
×
2589
          taosMemoryFree(tv);
×
2590
          return terrno;
×
2591
        }
2592
        if (oldTag) { taosMemoryFreeClear(oldTag->pData); taosMemoryFree(oldTag); }
13,132 ✔
2593
      }
2594
    } else {
2595
      SResolveWorkItem next = {0};
27,135 ✔
2596
      next.originVtbUid = w->originVtbUid;
27,135 ✔
2597
      next.originCid    = w->originCid;
27,135 ✔
2598
      next.kind         = r->nextRef.kind;
27,135 ✔
2599
      tstrncpy(next.refDbName,    r->nextRef.refDbName,    TSDB_DB_NAME_LEN);
27,135 ✔
2600
      tstrncpy(next.refTableName, r->nextRef.refTableName, TSDB_TABLE_NAME_LEN);
27,135 ✔
2601
      tstrncpy(next.refColName,   r->nextRef.refColName,   TSDB_COL_NAME_LEN);
27,135 ✔
2602
      if (taosArrayPush(nextWorkList, &next) == NULL) {
27,135 ✔
2603
        taosMemoryFreeClear(r->tagData);
×
2604
        return terrno;
×
2605
      }
2606
      stDebug("vgId:%d %s hop=%d uid=%" PRId64 " kind=%d cid=%d NEXT-HOP -> %s.%s.%s",
27,135 ✔
2607
              TD_VID(pVnode), __func__, hop, w->originVtbUid, next.kind, w->originCid,
2608
              next.refDbName, next.refTableName, next.refColName);
2609
      taosMemoryFreeClear(r->tagData);
27,135 ✔
2610
    }
2611
  }
2612
  return 0;
49,045 ✔
2613
}
2614

2615
// Debug-only dump of the final per-uid resolve result. Pure stDebug; called
2616
// once at the end of streamResolveVTableRefChain to ease post-mortem.
2617
static void streamLogFinalResolveResult(SVnode *pVnode, SSHashObj *uid2Result) {
46,054 ✔
2618
  void   *p  = NULL;
46,054 ✔
2619
  int32_t it = 0;
46,054 ✔
2620
  while ((p = tSimpleHashIterate(uid2Result, p, &it)) != NULL) {
114,203 ✔
2621
    int64_t               uid = *(int64_t *)tSimpleHashGetKey(p, NULL);
68,224 ✔
2622
    SVTableResolveResult *res = *(SVTableResolveResult **)p;
68,224 ✔
2623
    stDebug("vgId:%d %s FINAL uid=%" PRId64 " colMapSz=%d tagMapSz=%d",
68,224 ✔
2624
            TD_VID(pVnode), __func__, uid,
2625
            tSimpleHashGetSize(res->colMap),
2626
            tSimpleHashGetSize(res->tagMap));
2627
    void *cp = NULL; int32_t ci = 0;
68,224 ✔
2628
    while ((cp = tSimpleHashIterate(res->colMap, cp, &ci)) != NULL) {
204,330 ✔
2629
      col_id_t         cid  = *(col_id_t *)tSimpleHashGetKey(cp, NULL);
136,106 ✔
2630
      SColResolveItem *item = *(SColResolveItem **)cp;
136,106 ✔
2631
      stDebug("vgId:%d %s   FINAL uid=%" PRId64 " COL cid=%d hasRef=%d ref=%s.%s.%s",
136,106 ✔
2632
              TD_VID(pVnode), __func__, uid, cid, item ? item->hasRef : -1,
2633
              item ? item->refDbName : "", item ? item->refTableName : "",
2634
              item ? item->refColName : "");
2635
    }
2636
    void *tp = NULL; int32_t ti = 0;
68,224 ✔
2637
    while ((tp = tSimpleHashIterate(res->tagMap, tp, &ti)) != NULL) {
96,949 ✔
2638
      col_id_t   cid = *(col_id_t *)tSimpleHashGetKey(tp, NULL);
28,725 ✔
2639
      STagValue *tv  = *(STagValue **)tp;
28,725 ✔
2640
      stDebug("vgId:%d %s   FINAL uid=%" PRId64 " TAG cid=%d type=%d nLen=%d",
28,725 ✔
2641
              TD_VID(pVnode), __func__, uid, cid, tv ? tv->type : -1, tv ? tv->nLen : -1);
2642
    }
2643
  }
2644
}
46,054 ✔
2645

2646
int32_t streamResolveVTableRefChain(SVnode *pVnode, SStreamVTableInfoCache *pCache,
46,129 ✔
2647
                                    SStreamTriggerReaderInfo *pReaderInfo, int64_t ver,
2648
                                    SArray *vtbUids, SArray *virtColCids, SArray *virtTagCids,
2649
                                    SSHashObj **ppUid2Result) {
2650
  int32_t    code         = 0;
46,129 ✔
2651
  SArray    *workList     = NULL;
46,129 ✔
2652
  SArray    *nextWorkList = NULL;
46,129 ✔
2653
  SArray    *rspItems     = NULL;
46,129 ✔
2654
  SSHashObj *uid2Result   = NULL;
46,129 ✔
2655

2656
  SArray    *fullUids     = NULL;
46,129 ✔
2657
  SArray    *pTableListArray = NULL;
46,188 ✔
2658

2659
  if (pVnode == NULL || ppUid2Result == NULL) return TSDB_CODE_INVALID_PARA;
46,188 ✔
2660
  *ppUid2Result = NULL;
46,188 ✔
2661

2662
  // Invalidate per-table ref cache at the start of each full resolve cycle.
2663
  // The cache is only useful within a single multi-hop resolve call to avoid
2664
  // redundant RPC for the same (db,table,col) across hops; stale results from
2665
  // a previous cycle could mask schema changes.
2666
  if (pCache) {
46,188 ✔
2667
    taosHashClear(pCache->tblRefCache);
36,987 ✔
2668
  }
2669

2670
  stDebug("vgId:%d %s enter: ver=%" PRId64 " vtbUids=%d virtCols=%d virtTags=%d", TD_VID(pVnode),
46,188 ✔
2671
          __func__, ver, (int32_t)taosArrayGetSize(vtbUids),
2672
          (int32_t)taosArrayGetSize(virtColCids), (int32_t)taosArrayGetSize(virtTagCids));
2673

2674
  // Full-uid branch: pull live uids from the reader's table list.
2675
  if (vtbUids == NULL || taosArrayGetSize(vtbUids) == 0) {
46,188 ✔
2676
    code = streamCollectActiveVtableUids(pReaderInfo, &pTableListArray, &fullUids);
4,269 ✔
2677
    if (code != 0) goto _end;
4,269 ✔
2678
    vtbUids = fullUids;
4,269 ✔
2679
    stDebug("vgId:%d %s full-uid branch: tableList=%d activeUids=%d", TD_VID(pVnode), __func__,
4,269 ✔
2680
            (int32_t)taosArrayGetSize(pTableListArray), (int32_t)taosArrayGetSize(fullUids));
2681
  }
2682

2683
  uid2Result = tSimpleHashInit(64, taosGetDefaultHashFunction(TSDB_DATA_TYPE_BIGINT));
46,188 ✔
2684
  if (uid2Result == NULL) { code = terrno; goto _end; }
46,188 ✔
2685
  tSimpleHashSetFreeFp(uid2Result, streamVTableResolveResultDestroy);
46,188 ✔
2686

2687
  workList = taosArrayInit(64, sizeof(SResolveWorkItem));
46,188 ✔
2688
  if (workList == NULL) { code = terrno; goto _end; }
46,188 ✔
2689

2690
  // 1. seed work-list. H2 v0.5: streamPushInitialWorkItemsForUid swallows
2691
  //    top-level uid-not-exist (warn + return 0 without entry); any other
2692
  //    error (col/tag not in ref triple, OOM) is propagated upward so the
2693
  //    caller (reader -> trigger -> mnode) can fail-fast and trigger a
2694
  //    redeploy.
2695
  int32_t nUid = (int32_t)taosArrayGetSize(vtbUids);
46,188 ✔
2696
  for (int32_t i = 0; i < nUid; ++i) {
123,907 ✔
2697
    int64_t uid = *(int64_t *)taosArrayGet(vtbUids, i);
77,778 ✔
2698
    int32_t rc  = streamPushInitialWorkItemsForUid(pVnode, uid, virtColCids, virtTagCids, workList, uid2Result);
77,778 ✔
2699
    if (rc == 0) continue;
77,719 ✔
2700
    stError("vgId:%d %s seed uid=%" PRId64 " push rc=0x%x -> propagate (strict)",
×
2701
            TD_VID(pVnode), __func__, uid, rc);
2702
    code = rc;
×
2703
    goto _end;
×
2704
  }
2705
  stDebug("vgId:%d %s after seed: workListSz=%d uid2ResultSz=%d",
46,129 ✔
2706
          TD_VID(pVnode), __func__,
2707
          (int32_t)taosArrayGetSize(workList),
2708
          tSimpleHashGetSize(uid2Result));
2709

2710
  // 2. main hop loop
2711
  for (int32_t hop = 0; hop < STREAM_VTB_MAX_HOPS; ++hop) {
95,158 ✔
2712
    int32_t cur = (int32_t)taosArrayGetSize(workList);
95,158 ✔
2713
    stDebug("vgId:%d %s hop=%d workListSz=%d", TD_VID(pVnode), __func__, hop, cur);
95,158 ✔
2714
    if (cur == 0) break;
95,158 ✔
2715

2716
    rspItems = taosArrayInit(cur, sizeof(SVTableRefResolveRspItem));
49,179 ✔
2717
    if (rspItems == NULL) { code = terrno; goto _end; }
49,179 ✔
2718

2719
    int32_t rc = streamCallResolveBatched(pVnode, pCache, ver, workList, rspItems);
49,179 ✔
2720
    if (rc != 0) {
49,119 ✔
2721
      // H2 v0.5: any error (OOM, routing, RPC) propagates immediately.
2722
      code = rc;
×
2723
      goto _end;
×
2724
    }
2725

2726
    nextWorkList = taosArrayInit(cur, sizeof(SResolveWorkItem));
49,119 ✔
2727
    if (nextWorkList == NULL) { code = terrno; goto _end; }
49,179 ✔
2728

2729
    code = streamConsumeOneHopResults(pVnode, hop, workList, rspItems, uid2Result, nextWorkList);
49,179 ✔
2730
    if (code != 0) goto _end;
49,179 ✔
2731

2732
    taosArrayDestroy(rspItems); rspItems = NULL;
49,045 ✔
2733
    taosArrayDestroy(workList);
48,970 ✔
2734
    workList     = nextWorkList;
48,970 ✔
2735
    nextWorkList = NULL;
48,970 ✔
2736
  }
2737

2738
  // 3. hop overflow: any leftover work-items mean the chain exceeded MAX_HOPS.
2739
  //    H2 v0.5: report TSDB_CODE_STREAM_VTB_REF_TOO_DEEP rather than silently
2740
  //    skipping the offending uids.
2741
  if (workList != NULL) {
45,979 ✔
2742
    int32_t leftover = (int32_t)taosArrayGetSize(workList);
45,979 ✔
2743
    if (leftover > 0) {
45,979 ✔
2744
      for (int32_t i = 0; i < leftover; ++i) {
×
2745
        SResolveWorkItem *w = taosArrayGet(workList, i);
×
2746
        stError("vgId:%d %s OVERFLOW uid=%" PRId64 " kind=%d cid=%d ref=%s.%s.%s",
×
2747
                TD_VID(pVnode), __func__, w->originVtbUid, w->kind, w->originCid,
2748
                w->refDbName, w->refTableName, w->refColName);
2749
      }
2750
      stError("vgId:%d %s HOP_OVERFLOW leftover=%d -> propagate TOO_DEEP",
×
2751
              TD_VID(pVnode), __func__, leftover);
2752
      code = TSDB_CODE_STREAM_VTB_REF_TOO_DEEP;
×
2753
      goto _end;
×
2754
    }
2755
  }
2756

2757
  // Final dump: per-uid colMap/tagMap contents.
2758
  streamLogFinalResolveResult(pVnode, uid2Result);
45,979 ✔
2759

2760
  *ppUid2Result = uid2Result;
46,054 ✔
2761
  uid2Result    = NULL;
46,054 ✔
2762

2763
_end:
46,188 ✔
2764
  stDebug("vgId:%d %s exit: code=0x%x outUidCnt=%d", TD_VID(pVnode), __func__, code,
46,188 ✔
2765
          uid2Result ? tSimpleHashGetSize(uid2Result) :
2766
          (*ppUid2Result ? tSimpleHashGetSize(*ppUid2Result) : 0));
2767
  if (fullUids        != NULL) taosArrayDestroy(fullUids);
46,188 ✔
2768
  if (pTableListArray != NULL) taosArrayDestroyP(pTableListArray, taosMemFree);
46,188 ✔
2769
  if (workList     != NULL) taosArrayDestroy(workList);
46,188 ✔
2770
  if (nextWorkList != NULL) taosArrayDestroy(nextWorkList);
46,188 ✔
2771
  if (rspItems     != NULL) {
46,188 ✔
2772
    int32_t m = (int32_t)taosArrayGetSize(rspItems);
134 ✔
2773
    for (int32_t i = 0; i < m; ++i) {
268 ✔
2774
      SVTableRefResolveRspItem *r = taosArrayGet(rspItems, i);
134 ✔
2775
      taosMemoryFreeClear(r->tagData);
134 ✔
2776
    }
2777
    taosArrayDestroy(rspItems);
134 ✔
2778
  }
2779
  tSimpleHashCleanup(uid2Result);
46,188 ✔
2780

2781
  return code;
46,188 ✔
2782
}
2783

2784
// ============================================================================
2785
// C3: vchild-tag chain helper (executor side)
2786
// vnodeResolveVTableTagChain
2787
//
2788
// For trigger streams whose source is a virtual super table, executor's
2789
// `getColInfoResultForGroupbyForStream` needs literal tag values per vchild to
2790
// compute the partition groupId. The default `metaGetTableTagsByUidsVersion`
2791
// only reads ctbEntry.pTags directly, so col-ref tags resolve to NULL and all
2792
// vchildren collapse into the same group.
2793
//
2794
// This helper post-processes the STUidTagInfo list: when suid is a virtual
2795
// stable, each vchild uid is fed into `streamResolveVTableRefChain` (which
2796
// already handles multi-hop and cross-vnode resolution), and the returned tag
2797
// values are repacked into a fresh STag in stable-schemaTag order. Failures
2798
// per uid are best-effort and leave the original pTagVal untouched.
2799
// ============================================================================
2800
int32_t vnodeResolveVTableTagChain(void *pVnode, int64_t suid, SArray *pUidTagList) {
698,207 ✔
2801
  if (pVnode == NULL || pUidTagList == NULL) return 0;
698,207 ✔
2802

2803
  int32_t      code         = 0;
698,683 ✔
2804
  SVnode      *pVn          = (SVnode *)pVnode;
698,683 ✔
2805
  stTrace("vgId:%d %s ENTER suid=%" PRId64 " nUids=%d", TD_VID(pVn), __func__, suid,
698,683 ✔
2806
          (int32_t)taosArrayGetSize(pUidTagList));
2807
  SMetaReader  mr           = {0};
698,683 ✔
2808
  bool         readerInited = false;
698,821 ✔
2809
  SArray      *uids         = NULL;
698,821 ✔
2810
  SArray      *tagCids      = NULL;
698,821 ✔
2811
  SArray      *tagVals      = NULL;
698,821 ✔
2812
  SSHashObj   *uid2Result = NULL;
698,821 ✔
2813
  int32_t      nTagCols     = 0;
698,407 ✔
2814

2815
  int32_t nUids = (int32_t)taosArrayGetSize(pUidTagList);
698,407 ✔
2816
  if (nUids == 0) return 0;
698,407 ✔
2817

2818
  // 1) confirm suid refers to a virtual super table; otherwise no-op.
2819
  metaReaderDoInit(&mr, pVn->pMeta, META_READER_LOCK, 0);
698,407 ✔
2820
  readerInited = true;
698,821 ✔
2821
  if (metaReaderGetTableEntryByUid(&mr, suid) != 0) {
698,821 ✔
2822
    stDebug("vgId:%d %s metaReader miss suid=%" PRId64, TD_VID(pVn), __func__, suid);
555 ✔
2823
    goto _end;
555 ✔
2824
  }
2825
  if (mr.me.type != TSDB_SUPER_TABLE || !TABLE_IS_VIRTUAL(mr.me.flags)) {
697,820 ✔
2826
    stDebug("vgId:%d %s skip suid=%" PRId64 " type=%d flags=0x%x", TD_VID(pVn), __func__,
696,040 ✔
2827
            suid, (int32_t)mr.me.type, (uint32_t)mr.me.flags);
2828
    goto _end;
696,373 ✔
2829
  }
2830
  SSchema *pTagSchema = mr.me.stbEntry.schemaTag.pSchema;
1,780 ✔
2831
  nTagCols = mr.me.stbEntry.schemaTag.nCols;
1,780 ✔
2832
  if (pTagSchema == NULL || nTagCols <= 0) {
1,780 ✔
2833
    goto _end;
×
2834
  }
2835

2836
  stTrace("vgId:%d %s suid=%" PRId64 " nUids=%d nTagCols=%d", TD_VID(pVn), __func__,
1,780 ✔
2837
          suid, nUids, nTagCols);
2838

2839
  tagCids = taosArrayInit(nTagCols, sizeof(col_id_t));
1,780 ✔
2840
  if (tagCids == NULL) {
1,780 ✔
2841
    code = terrno;
×
2842
    goto _end;
×
2843
  }
2844
  for (int32_t i = 0; i < nTagCols; ++i) {
4,013 ✔
2845
    col_id_t cid = pTagSchema[i].colId;
2,115 ✔
2846
    if (taosArrayPush(tagCids, &cid) == NULL) {
2,115 ✔
2847
      code = terrno;
×
2848
      goto _end;
×
2849
    }
2850
  }
2851

2852
  metaReaderClear(&mr);
1,898 ✔
2853
  readerInited = false;
1,839 ✔
2854

2855
  // 2) build uid array for the chain resolver.
2856
  uids = taosArrayInit(nUids, sizeof(int64_t));
1,839 ✔
2857
  if (uids == NULL) { code = terrno; goto _end; }
1,780 ✔
2858
  for (int32_t i = 0; i < nUids; ++i) {
5,514 ✔
2859
    STUidTagInfo *p = taosArrayGet(pUidTagList, i);
3,675 ✔
2860
    if (p == NULL) continue;
3,616 ✔
2861
    int64_t uid = (int64_t)p->uid;
3,616 ✔
2862
    if (taosArrayPush(uids, &uid) == NULL) { code = terrno; goto _end; }
3,734 ✔
2863
  }
2864

2865
  // 3) chain-resolve.
2866
  code = streamResolveVTableRefChain(pVn, NULL, NULL, -1, uids, NULL, tagCids, &uid2Result);
1,839 ✔
2867
  if (code != 0 || uid2Result == NULL) {
1,839 ✔
2868
    stTrace("vgId:%d %s chain resolve rc=0x%x uid2Result=%p", TD_VID(pVn), __func__,
×
2869
            code, (void *)uid2Result);
2870
    code = 0;  // best-effort; do not propagate to caller
×
2871
    goto _end;
×
2872
  }
2873

2874
  // 4) rebuild STag per uid by merging:
2875
  //      - literal STagVals already present in p->pTagVal (vchild may declare
2876
  //        some tags as plain literals and only some as colRefs)
2877
  //      - chain-resolved STagValues from streamResolveVTableRefChain
2878
  //    The chain resolver only fills cids that appeared as colRefs, so without
2879
  //    this merge the literal tags would be lost when we tTagNew a fresh STag.
2880
  for (int32_t i = 0; i < nUids; ++i) {
5,573 ✔
2881
    STUidTagInfo *p = taosArrayGet(pUidTagList, i);
3,734 ✔
2882
    if (p == NULL) continue;
3,734 ✔
2883
    int64_t uid = (int64_t)p->uid;
3,734 ✔
2884

2885
    SVTableResolveResult **ppRes =
×
2886
        (SVTableResolveResult **)tSimpleHashGet(uid2Result, &uid, sizeof(uid));
3,734 ✔
2887
    SVTableResolveResult  *pRes  = (ppRes != NULL) ? *ppRes : NULL;
3,734 ✔
2888
    bool hasResolvedTags = (pRes != NULL && pRes->tagMap != NULL && tSimpleHashGetSize(pRes->tagMap) > 0);
3,734 ✔
2889
    stTrace("vgId:%d %s merge uid=%" PRId64 " ppRes=%p pRes=%p tagMapSz=%d",
3,734 ✔
2890
            TD_VID(pVn), __func__, uid, (void *)ppRes, (void *)pRes,
2891
            (pRes && pRes->tagMap) ? tSimpleHashGetSize(pRes->tagMap) : -1);
2892

2893
    if (tagVals != NULL) {
3,734 ✔
2894
      taosArrayClear(tagVals);
1,895 ✔
2895
    } else {
2896
      tagVals = taosArrayInit(nTagCols, sizeof(STagVal));
1,839 ✔
2897
      if (tagVals == NULL) { code = terrno; goto _end; }
1,839 ✔
2898
    }
2899

2900
    bool anyChange = false;
3,734 ✔
2901
    for (int32_t j = 0; j < nTagCols; ++j) {
8,406 ✔
2902
      col_id_t cid = *(col_id_t *)taosArrayGet(tagCids, j);
4,672 ✔
2903

2904
      // Prefer chain-resolved value when present (overrides any stale literal).
2905
      if (hasResolvedTags) {
4,672 ✔
2906
        STagValue **ppTV = (STagValue **)tSimpleHashGet(pRes->tagMap, &cid, sizeof(cid));
4,672 ✔
2907
        if (ppTV != NULL && *ppTV != NULL) {
4,672 ✔
2908
          STagValue *tv = *ppTV;
4,672 ✔
2909
          if (tv->pData != NULL && tv->nLen > 0) {
4,672 ✔
2910
            STagVal v = {0};
4,672 ✔
2911
            v.cid  = cid;
4,672 ✔
2912
            v.type = tv->type;
4,672 ✔
2913
            if (IS_VAR_DATA_TYPE(tv->type)) {
4,672 ✔
2914
              v.nData = (uint32_t)tv->nLen;
165 ✔
2915
              v.pData = (uint8_t *)tv->pData;
165 ✔
2916
            } else {
2917
              int32_t copyLen = tv->nLen < (int32_t)sizeof(int64_t) ? tv->nLen : (int32_t)sizeof(int64_t);
4,507 ✔
2918
              memcpy(&v.i64, tv->pData, copyLen);
4,507 ✔
2919
            }
2920
            if (taosArrayPush(tagVals, &v) == NULL) { code = terrno; goto _end; }
4,672 ✔
2921
            anyChange = true;
4,672 ✔
2922
            continue;
4,672 ✔
2923
          }
2924
        }
2925
      }
2926

2927
      // Fall back to the literal tag in the original STag, if any.
2928
      if (p->pTagVal != NULL) {
×
2929
        STagVal probe = {.cid = cid};
×
2930
        if (tTagGet((const STag *)p->pTagVal, &probe)) {
×
2931
          if (taosArrayPush(tagVals, &probe) == NULL) { code = terrno; goto _end; }
×
2932
        }
2933
      }
2934
    }
2935

2936
    if (!anyChange) continue;  // nothing was resolved -> keep the original STag
3,734 ✔
2937

2938
    STag   *pNewTag = NULL;
3,734 ✔
2939
    int32_t rc      = tTagNew(tagVals, 1, false, &pNewTag);
3,734 ✔
2940
    if (rc != 0 || pNewTag == NULL) {
3,734 ✔
2941
      stDebug("vgId:%d %s uid=%" PRId64 " tTagNew rc=0x%x -> keep original", TD_VID(pVn),
×
2942
              __func__, uid, rc);
2943
      continue;
×
2944
    }
2945
    if (p->pTagVal != NULL) taosMemoryFree(p->pTagVal);
3,734 ✔
2946
    p->pTagVal = pNewTag;
3,734 ✔
2947
    stDebug("vgId:%d %s uid=%" PRId64 " rebuilt STag with %d tag(s) (literals merged)",
3,734 ✔
2948
            TD_VID(pVn), __func__, uid, (int32_t)taosArrayGetSize(tagVals));
2949
  }
2950

2951
_end:
698,704 ✔
2952
  if (readerInited) metaReaderClear(&mr);
698,825 ✔
2953
  if (uids    != NULL) taosArrayDestroy(uids);
698,821 ✔
2954
  if (tagCids != NULL) taosArrayDestroy(tagCids);
698,821 ✔
2955
  if (tagVals != NULL) taosArrayDestroy(tagVals);
698,821 ✔
2956
  tSimpleHashCleanup(uid2Result);
698,821 ✔
2957
  return code;
698,677 ✔
2958
}
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