indexTfile.c 29.1 KB
Newer Older
dengyihao's avatar
dengyihao 已提交
1 2
/*
 * Copyright (c) 2019 TAOS Data, Inc. <jhtao@taosdata.com>
dengyihao's avatar
dengyihao 已提交
3
p *
dengyihao's avatar
dengyihao 已提交
4 5 6 7 8 9 10 11 12 13 14 15
 * This program is free software: you can use, redistribute, and/or modify
 * it under the terms of the GNU Affero General Public License, version 3
 * or later ("AGPL"), as published by the Free Software Foundation.
 *
 * This program is distributed in the hope that it will be useful, but WITHOUT
 * ANY WARRANTY; without even the implied warranty of MERCHANTABILITY or
 * FITNESS FOR A PARTICULAR PURPOSE.
 *
 * You should have received a copy of the GNU Affero General Public License
 * along with this program. If not, see <http://www.gnu.org/licenses/>.
 */

dengyihao's avatar
dengyihao 已提交
16
#include "indexTfile.h"
17
#include "index.h"
dengyihao's avatar
dengyihao 已提交
18 19 20 21
#include "indexComm.h"
#include "indexFst.h"
#include "indexFstCountingWriter.h"
#include "indexUtil.h"
dengyihao's avatar
dengyihao 已提交
22
#include "taosdef.h"
dengyihao's avatar
dengyihao 已提交
23
#include "tcoding.h"
dengyihao's avatar
dengyihao 已提交
24
#include "tcompare.h"
dengyihao's avatar
dengyihao 已提交
25

dengyihao's avatar
dengyihao 已提交
26 27
const static uint64_t tfileMagicNumber = 0xdb4775248b80fb57ull;

dengyihao's avatar
dengyihao 已提交
28 29 30 31 32 33 34
typedef struct TFileFstIter {
  FstStreamBuilder* fb;
  StreamWithState*  st;
  AutomationCtx*    ctx;
  TFileReader*      rdr;
} TFileFstIter;

dengyihao's avatar
dengyihao 已提交
35 36
#define TF_TABLE_TATOAL_SIZE(sz) (sizeof(sz) + sz * sizeof(uint64_t))

dengyihao's avatar
dengyihao 已提交
37
static int  tfileUidCompare(const void* a, const void* b);
dengyihao's avatar
dengyihao 已提交
38
static int  tfileStrCompare(const void* a, const void* b);
dengyihao's avatar
dengyihao 已提交
39 40
static int  tfileValueCompare(const void* a, const void* b, const void* param);
static void tfileSerialTableIdsToBuf(char* buf, SArray* tableIds);
dengyihao's avatar
dengyihao 已提交
41

dengyihao's avatar
dengyihao 已提交
42
static int tfileWriteHeader(TFileWriter* writer);
dengyihao's avatar
dengyihao 已提交
43
static int tfileWriteFstOffset(TFileWriter* tw, int32_t offset);
dengyihao's avatar
dengyihao 已提交
44
static int tfileWriteData(TFileWriter* write, TFileValue* tval);
dengyihao's avatar
dengyihao 已提交
45
static int tfileWriteFooter(TFileWriter* write);
dengyihao's avatar
dengyihao 已提交
46

dengyihao's avatar
dengyihao 已提交
47
// handle file corrupt later
dengyihao's avatar
dengyihao 已提交
48 49
static int tfileReaderLoadHeader(TFileReader* reader);
static int tfileReaderLoadFst(TFileReader* reader);
dengyihao's avatar
dengyihao 已提交
50
static int tfileReaderVerify(TFileReader* reader);
dengyihao's avatar
dengyihao 已提交
51
static int tfileReaderLoadTableIds(TFileReader* reader, int32_t offset, SArray* result);
dengyihao's avatar
dengyihao 已提交
52

dengyihao's avatar
dengyihao 已提交
53 54 55 56 57 58 59
static SArray* tfileGetFileList(const char* path);
static int     tfileRmExpireFile(SArray* result);
static void    tfileDestroyFileName(void* elem);
static int     tfileCompare(const void* a, const void* b);
static int     tfileParseFileName(const char* filename, uint64_t* suid, char* col, int* version);
static void    tfileGenFileName(char* filename, uint64_t suid, const char* col, int version);
static void    tfileGenFileFullName(char* fullname, const char* path, uint64_t suid, const char* col, int32_t version);
dengyihao's avatar
dengyihao 已提交
60 61 62 63 64 65 66
/*
 * search from  tfile
 */
static int32_t tfSearchTerm(void* reader, SIndexTerm* tem, SIdxTempResult* tr);
static int32_t tfSearchPrefix(void* reader, SIndexTerm* tem, SIdxTempResult* tr);
static int32_t tfSearchSuffix(void* reader, SIndexTerm* tem, SIdxTempResult* tr);
static int32_t tfSearchRegex(void* reader, SIndexTerm* tem, SIdxTempResult* tr);
dengyihao's avatar
dengyihao 已提交
67 68 69 70
static int32_t tfSearchLessThan(void* reader, SIndexTerm* tem, SIdxTempResult* tr);
static int32_t tfSearchLessEqual(void* reader, SIndexTerm* tem, SIdxTempResult* tr);
static int32_t tfSearchGreaterThan(void* reader, SIndexTerm* tem, SIdxTempResult* tr);
static int32_t tfSearchGreaterEqual(void* reader, SIndexTerm* tem, SIdxTempResult* tr);
dengyihao's avatar
dengyihao 已提交
71 72
static int32_t tfSearchRange(void* reader, SIndexTerm* tem, SIdxTempResult* tr);

dengyihao's avatar
dengyihao 已提交
73 74
static int32_t tfSearchCompareFunc(void* reader, SIndexTerm* tem, SIdxTempResult* tr, RangeType ctype);

dengyihao's avatar
dengyihao 已提交
75
static int32_t (*tfSearch[])(void* reader, SIndexTerm* tem, SIdxTempResult* tr) = {
dengyihao's avatar
dengyihao 已提交
76 77
    tfSearchTerm,      tfSearchPrefix,      tfSearchSuffix,       tfSearchRegex, tfSearchLessThan,
    tfSearchLessEqual, tfSearchGreaterThan, tfSearchGreaterEqual, tfSearchRange};
dengyihao's avatar
dengyihao 已提交
78

dengyihao's avatar
dengyihao 已提交
79
TFileCache* tfileCacheCreate(const char* path) {
wafwerar's avatar
wafwerar 已提交
80
  TFileCache* tcache = taosMemoryCalloc(1, sizeof(TFileCache));
dengyihao's avatar
dengyihao 已提交
81 82 83
  if (tcache == NULL) {
    return NULL;
  }
84 85 86 87

  tcache->tableCache = taosHashInit(8, taosGetDefaultHashFunction(TSDB_DATA_TYPE_BINARY), true, HASH_ENTRY_LOCK);
  tcache->capacity = 64;

dengyihao's avatar
dengyihao 已提交
88
  SArray* files = tfileGetFileList(path);
dengyihao's avatar
dengyihao 已提交
89
  for (size_t i = 0; i < taosArrayGetSize(files); i++) {
dengyihao's avatar
dengyihao 已提交
90
    char* file = taosArrayGetP(files, i);
dengyihao's avatar
dengyihao 已提交
91

dengyihao's avatar
dengyihao 已提交
92
    WriterCtx* wc = writerCtxCreate(TFile, file, true, 1024 * 1024 * 64);
93
    if (wc == NULL) {
dengyihao's avatar
dengyihao 已提交
94
      indexError("failed to open index:%s", file);
95
      goto End;
dengyihao's avatar
dengyihao 已提交
96
    }
dengyihao's avatar
dengyihao 已提交
97

dengyihao's avatar
dengyihao 已提交
98
    TFileReader* reader = tfileReaderCreate(wc);
dengyihao's avatar
dengyihao 已提交
99 100 101 102
    if (reader == NULL) {
      indexInfo("skip invalid file: %s", file);
      continue;
    }
dengyihao's avatar
dengyihao 已提交
103
    TFileHeader* header = &reader->header;
dengyihao's avatar
dengyihao 已提交
104
    ICacheKey    key = {.suid = header->suid, .colName = header->colName, .nColName = strlen(header->colName)};
dengyihao's avatar
dengyihao 已提交
105

dengyihao's avatar
dengyihao 已提交
106
    char    buf[128] = {0};
dengyihao's avatar
dengyihao 已提交
107 108 109
    int32_t sz = indexSerialCacheKey(&key, buf);
    assert(sz < sizeof(buf));
    taosHashPut(tcache->tableCache, buf, sz, &reader, sizeof(void*));
dengyihao's avatar
dengyihao 已提交
110
    tfileReaderRef(reader);
dengyihao's avatar
dengyihao 已提交
111
  }
dengyihao's avatar
dengyihao 已提交
112
  taosArrayDestroyEx(files, tfileDestroyFileName);
dengyihao's avatar
dengyihao 已提交
113
  return tcache;
dengyihao's avatar
dengyihao 已提交
114
End:
115
  tfileCacheDestroy(tcache);
dengyihao's avatar
dengyihao 已提交
116
  taosArrayDestroyEx(files, tfileDestroyFileName);
117
  return NULL;
dengyihao's avatar
dengyihao 已提交
118
}
dengyihao's avatar
dengyihao 已提交
119
void tfileCacheDestroy(TFileCache* tcache) {
dengyihao's avatar
dengyihao 已提交
120 121 122
  if (tcache == NULL) {
    return;
  }
123
  // free table cache
dengyihao's avatar
dengyihao 已提交
124
  TFileReader** reader = taosHashIterate(tcache->tableCache, NULL);
125
  while (reader) {
dengyihao's avatar
dengyihao 已提交
126
    TFileReader* p = *reader;
127 128
    indexInfo("drop table cache suid: %" PRIu64 ", colName: %s, colType: %d", p->header.suid, p->header.colName,
              p->header.colType);
dengyihao's avatar
dengyihao 已提交
129

dengyihao's avatar
dengyihao 已提交
130
    tfileReaderUnRef(p);
131 132 133
    reader = taosHashIterate(tcache->tableCache, reader);
  }
  taosHashCleanup(tcache->tableCache);
wafwerar's avatar
wafwerar 已提交
134
  taosMemoryFree(tcache);
dengyihao's avatar
dengyihao 已提交
135 136
}

dengyihao's avatar
dengyihao 已提交
137 138 139 140 141
TFileReader* tfileCacheGet(TFileCache* tcache, ICacheKey* key) {
  char    buf[128] = {0};
  int32_t sz = indexSerialCacheKey(key, buf);
  assert(sz < sizeof(buf));
  TFileReader** reader = taosHashGet(tcache->tableCache, buf, sz);
dengyihao's avatar
dengyihao 已提交
142 143 144
  if (reader == NULL) {
    return NULL;
  }
145
  tfileReaderRef(*reader);
dengyihao's avatar
dengyihao 已提交
146

147
  return *reader;
dengyihao's avatar
dengyihao 已提交
148
}
dengyihao's avatar
dengyihao 已提交
149 150 151
void tfileCachePut(TFileCache* tcache, ICacheKey* key, TFileReader* reader) {
  char    buf[128] = {0};
  int32_t sz = indexSerialCacheKey(key, buf);
dengyihao's avatar
dengyihao 已提交
152
  // remove last version index reader
dengyihao's avatar
dengyihao 已提交
153
  TFileReader** p = taosHashGet(tcache->tableCache, buf, sz);
dengyihao's avatar
dengyihao 已提交
154
  if (p != NULL) {
dengyihao's avatar
dengyihao 已提交
155
    TFileReader* oldReader = *p;
dengyihao's avatar
dengyihao 已提交
156
    taosHashRemove(tcache->tableCache, buf, sz);
dengyihao's avatar
dengyihao 已提交
157
    oldReader->remove = true;
dengyihao's avatar
dengyihao 已提交
158
    tfileReaderUnRef(oldReader);
dengyihao's avatar
dengyihao 已提交
159
  }
dengyihao's avatar
dengyihao 已提交
160

dengyihao's avatar
dengyihao 已提交
161
  taosHashPut(tcache->tableCache, buf, sz, &reader, sizeof(void*));
dengyihao's avatar
dengyihao 已提交
162
  tfileReaderRef(reader);
dengyihao's avatar
dengyihao 已提交
163
  return;
164
}
dengyihao's avatar
dengyihao 已提交
165
TFileReader* tfileReaderCreate(WriterCtx* ctx) {
wafwerar's avatar
wafwerar 已提交
166
  TFileReader* reader = taosMemoryCalloc(1, sizeof(TFileReader));
dengyihao's avatar
dengyihao 已提交
167 168 169
  if (reader == NULL) {
    return NULL;
  }
dengyihao's avatar
dengyihao 已提交
170 171

  reader->ctx = ctx;
dengyihao's avatar
dengyihao 已提交
172 173

  if (0 != tfileReaderVerify(reader)) {
dengyihao's avatar
dengyihao 已提交
174
    indexError("invalid tfile, suid: %" PRIu64 ", colName: %s", reader->header.suid, reader->header.colName);
dengyihao's avatar
dengyihao 已提交
175
    tfileReaderDestroy(reader);
dengyihao's avatar
dengyihao 已提交
176 177 178
    return NULL;
  }
  // T_REF_INC(reader);
dengyihao's avatar
dengyihao 已提交
179
  if (0 != tfileReaderLoadHeader(reader)) {
180 181
    indexError("failed to load index header, suid: %" PRIu64 ", colName: %s", reader->header.suid,
               reader->header.colName);
dengyihao's avatar
dengyihao 已提交
182
    tfileReaderDestroy(reader);
dengyihao's avatar
dengyihao 已提交
183 184 185 186
    return NULL;
  }

  if (0 != tfileReaderLoadFst(reader)) {
dengyihao's avatar
dengyihao 已提交
187 188
    indexError("failed to load index fst, suid: %" PRIu64 ", colName: %s, errno: %d", reader->header.suid,
               reader->header.colName, errno);
dengyihao's avatar
dengyihao 已提交
189 190 191 192
    tfileReaderDestroy(reader);
    return NULL;
  }

193
  return reader;
dengyihao's avatar
dengyihao 已提交
194
}
dengyihao's avatar
dengyihao 已提交
195
void tfileReaderDestroy(TFileReader* reader) {
dengyihao's avatar
dengyihao 已提交
196 197 198
  if (reader == NULL) {
    return;
  }
199
  // T_REF_INC(reader);
dengyihao's avatar
dengyihao 已提交
200
  fstDestroy(reader->fst);
dengyihao's avatar
dengyihao 已提交
201
  writerCtxDestroy(reader->ctx, reader->remove);
wafwerar's avatar
wafwerar 已提交
202
  taosMemoryFree(reader);
dengyihao's avatar
dengyihao 已提交
203
}
dengyihao's avatar
dengyihao 已提交
204 205 206 207 208 209 210 211 212 213 214 215 216 217 218 219 220
static int32_t tfSearchTerm(void* reader, SIndexTerm* tem, SIdxTempResult* tr) {
  bool     hasJson = INDEX_TYPE_CONTAIN_EXTERN_TYPE(tem->colType, TSDB_DATA_TYPE_JSON);
  int      ret = 0;
  char*    p = tem->colVal;
  uint64_t sz = tem->nColVal;
  if (hasJson) {
    p = indexPackJsonData(tem);
    sz = strlen(p);
  }
  int64_t  st = taosGetTimestampUs();
  FstSlice key = fstSliceCreate(p, sz);
  uint64_t offset;
  if (fstGet(((TFileReader*)reader)->fst, &key, &offset)) {
    int64_t et = taosGetTimestampUs();
    int64_t cost = et - st;
    indexInfo("index: %" PRIu64 ", col: %s, colVal: %s, found table info in tindex, time cost: %" PRIu64 "us",
              tem->suid, tem->colName, tem->colVal, cost);
dengyihao's avatar
dengyihao 已提交
221

dengyihao's avatar
dengyihao 已提交
222 223 224 225 226 227 228 229 230 231 232
    ret = tfileReaderLoadTableIds((TFileReader*)reader, offset, tr->total);
    cost = taosGetTimestampUs() - et;
    indexInfo("index: %" PRIu64 ", col: %s, colVal: %s, load all table info, time cost: %" PRIu64 "us", tem->suid,
              tem->colName, tem->colVal, cost);
  }
  if (hasJson) {
    taosMemoryFree(p);
  }
  fstSliceDestroy(&key);
  return 0;
}
dengyihao's avatar
dengyihao 已提交
233

dengyihao's avatar
dengyihao 已提交
234 235 236 237 238 239 240 241 242 243 244 245 246 247 248 249 250 251 252 253 254 255 256 257 258 259 260 261 262
static int32_t tfSearchPrefix(void* reader, SIndexTerm* tem, SIdxTempResult* tr) {
  bool     hasJson = INDEX_TYPE_CONTAIN_EXTERN_TYPE(tem->colType, TSDB_DATA_TYPE_JSON);
  char*    p = tem->colVal;
  uint64_t sz = tem->nColVal;
  if (hasJson) {
    p = indexPackJsonData(tem);
    sz = strlen(p);
  }

  SArray* offsets = taosArrayInit(16, sizeof(uint64_t));

  AutomationCtx*         ctx = automCtxCreate((void*)p, AUTOMATION_PREFIX);
  FstStreamBuilder*      sb = fstSearch(((TFileReader*)reader)->fst, ctx);
  StreamWithState*       st = streamBuilderIntoStream(sb);
  StreamWithStateResult* rt = NULL;
  while ((rt = streamWithStateNextWith(st, NULL)) != NULL) {
    taosArrayPush(offsets, &(rt->out.out));
    swsResultDestroy(rt);
  }
  streamWithStateDestroy(st);
  fstStreamBuilderDestroy(sb);

  int32_t ret = 0;
  for (int i = 0; i < taosArrayGetSize(offsets); i++) {
    uint64_t offset = *(uint64_t*)taosArrayGet(offsets, i);
    ret = tfileReaderLoadTableIds((TFileReader*)reader, offset, tr->total);
    if (ret != 0) {
      indexError("failed to find target tablelist");
      return TSDB_CODE_TDB_FILE_CORRUPTED;
dengyihao's avatar
add UT  
dengyihao 已提交
263
    }
dengyihao's avatar
dengyihao 已提交
264
  }
dengyihao's avatar
dengyihao 已提交
265 266 267 268 269 270 271 272 273 274 275 276 277 278 279 280 281 282 283 284 285 286 287 288 289 290 291 292 293 294 295 296 297 298 299 300 301 302 303 304 305 306 307 308
  if (hasJson) {
    taosMemoryFree(p);
  }
  return 0;
}
static int32_t tfSearchSuffix(void* reader, SIndexTerm* tem, SIdxTempResult* tr) {
  bool hasJson = INDEX_TYPE_CONTAIN_EXTERN_TYPE(tem->colType, TSDB_DATA_TYPE_JSON);

  int      ret = 0;
  char*    p = tem->colVal;
  uint64_t sz = tem->nColVal;
  if (hasJson) {
    p = indexPackJsonData(tem);
    sz = strlen(p);
  }
  int64_t  st = taosGetTimestampUs();
  FstSlice key = fstSliceCreate(p, sz);
  /*impl later*/
  if (hasJson) {
    taosMemoryFree(p);
  }
  fstSliceDestroy(&key);
  return 0;
}
static int32_t tfSearchRegex(void* reader, SIndexTerm* tem, SIdxTempResult* tr) {
  bool hasJson = INDEX_TYPE_CONTAIN_EXTERN_TYPE(tem->colType, TSDB_DATA_TYPE_JSON);

  int      ret = 0;
  char*    p = tem->colVal;
  uint64_t sz = tem->nColVal;
  if (hasJson) {
    p = indexPackJsonData(tem);
    sz = strlen(p);
  }
  int64_t  st = taosGetTimestampUs();
  FstSlice key = fstSliceCreate(p, sz);
  /*impl later*/

  if (hasJson) {
    taosMemoryFree(p);
  }
  fstSliceDestroy(&key);
  return 0;
}
dengyihao's avatar
dengyihao 已提交
309 310 311 312 313 314 315 316 317 318 319 320 321 322 323 324 325 326 327 328 329 330 331 332 333 334 335 336 337

static int32_t tfSearchCompareFunc(void* reader, SIndexTerm* tem, SIdxTempResult* tr, RangeType type) {
  bool     hasJson = INDEX_TYPE_CONTAIN_EXTERN_TYPE(tem->colType, TSDB_DATA_TYPE_JSON);
  int      ret = 0;
  char*    p = tem->colVal;
  uint64_t sz = tem->nColVal;
  if (hasJson) {
    p = indexPackJsonData(tem);
    sz = strlen(p);
  }
  SArray* offsets = taosArrayInit(16, sizeof(uint64_t));

  AutomationCtx*    ctx = automCtxCreate((void*)p, AUTOMATION_ALWAYS);
  FstStreamBuilder* sb = fstSearch(((TFileReader*)reader)->fst, ctx);

  FstSlice h = fstSliceCreate((uint8_t*)p, sz);
  fstStreamBuilderSetRange(sb, &h, type);
  fstSliceDestroy(&h);

  StreamWithState*       st = streamBuilderIntoStream(sb);
  StreamWithStateResult* rt = NULL;
  while ((rt = streamWithStateNextWith(st, NULL)) != NULL) {
    taosArrayPush(offsets, &(rt->out.out));
    swsResultDestroy(rt);
  }
  streamWithStateDestroy(st);
  fstStreamBuilderDestroy(sb);
  return TSDB_CODE_SUCCESS;
}
dengyihao's avatar
dengyihao 已提交
338
static int32_t tfSearchLessThan(void* reader, SIndexTerm* tem, SIdxTempResult* tr) {
dengyihao's avatar
dengyihao 已提交
339
  return tfSearchCompareFunc(reader, tem, tr, LT);
dengyihao's avatar
dengyihao 已提交
340 341
}
static int32_t tfSearchLessEqual(void* reader, SIndexTerm* tem, SIdxTempResult* tr) {
dengyihao's avatar
dengyihao 已提交
342
  return tfSearchCompareFunc(reader, tem, tr, LE);
dengyihao's avatar
dengyihao 已提交
343 344
}
static int32_t tfSearchGreaterThan(void* reader, SIndexTerm* tem, SIdxTempResult* tr) {
dengyihao's avatar
dengyihao 已提交
345
  return tfSearchCompareFunc(reader, tem, tr, GT);
dengyihao's avatar
dengyihao 已提交
346 347
}
static int32_t tfSearchGreaterEqual(void* reader, SIndexTerm* tem, SIdxTempResult* tr) {
dengyihao's avatar
dengyihao 已提交
348
  return tfSearchCompareFunc(reader, tem, tr, GE);
dengyihao's avatar
dengyihao 已提交
349
}
dengyihao's avatar
dengyihao 已提交
350 351 352 353 354 355 356 357 358 359 360 361 362 363 364 365 366
static int32_t tfSearchRange(void* reader, SIndexTerm* tem, SIdxTempResult* tr) {
  bool     hasJson = INDEX_TYPE_CONTAIN_EXTERN_TYPE(tem->colType, TSDB_DATA_TYPE_JSON);
  int      ret = 0;
  char*    p = tem->colVal;
  uint64_t sz = tem->nColVal;
  if (hasJson) {
    p = indexPackJsonData(tem);
    sz = strlen(p);
  }
  int64_t  st = taosGetTimestampUs();
  FstSlice key = fstSliceCreate(p, sz);
  // uint64_t offset;
  // if (fstGet(((TFileReader*)reader)->fst, &key, &offset)) {
  //  int64_t et = taosGetTimestampUs();
  //  int64_t cost = et - st;
  //  indexInfo("index: %" PRIu64 ", col: %s, colVal: %s, found table info in tindex, time cost: %" PRIu64 "us",
  //            tem->suid, tem->colName, tem->colVal, cost);
dengyihao's avatar
dengyihao 已提交
367

dengyihao's avatar
dengyihao 已提交
368 369 370 371 372 373 374 375 376 377 378
  //  ret = tfileReaderLoadTableIds((TFileReader*)reader, offset, tr->total);
  //  cost = taosGetTimestampUs() - et;
  //  indexInfo("index: %" PRIu64 ", col: %s, colVal: %s, load all table info, time cost: %" PRIu64 "us", tem->suid,
  //            tem->colName, tem->colVal, cost);
  //}
  if (hasJson) {
    taosMemoryFree(p);
  }
  fstSliceDestroy(&key);
  return 0;
}
dengyihao's avatar
dengyihao 已提交
379

dengyihao's avatar
dengyihao 已提交
380 381 382 383 384 385 386 387 388 389 390 391
int tfileReaderSearch(TFileReader* reader, SIndexTermQuery* query, SIdxTempResult* tr) {
  SIndexTerm*     term = query->term;
  EIndexQueryType qtype = query->qType;
  if (qtype >= sizeof(tfSearch) / sizeof(tfSearch[0])) {
    indexInfo("index: %" PRIu64 ", col: %s, colVal: %s, not found table info in tindex", term->suid, term->colName,
              term->colVal);
    return -1;
  } else {
    return tfSearch[qtype](reader, term, tr);
  }
  tfileReaderUnRef(reader);
  return 0;
dengyihao's avatar
dengyihao 已提交
392 393
}

dengyihao's avatar
dengyihao 已提交
394 395
TFileWriter* tfileWriterOpen(char* path, uint64_t suid, int32_t version, const char* colName, uint8_t colType) {
  char fullname[256] = {0};
dengyihao's avatar
dengyihao 已提交
396
  tfileGenFileFullName(fullname, path, suid, colName, version);
dengyihao's avatar
dengyihao 已提交
397
  // indexInfo("open write file name %s", fullname);
dengyihao's avatar
dengyihao 已提交
398
  WriterCtx* wcx = writerCtxCreate(TFile, fullname, false, 1024 * 1024 * 64);
dengyihao's avatar
dengyihao 已提交
399 400 401
  if (wcx == NULL) {
    return NULL;
  }
dengyihao's avatar
dengyihao 已提交
402 403 404 405 406 407 408 409 410

  TFileHeader tfh = {0};
  tfh.suid = suid;
  tfh.version = version;
  memcpy(tfh.colName, colName, strlen(colName));
  tfh.colType = colType;

  return tfileWriterCreate(wcx, &tfh);
}
dengyihao's avatar
dengyihao 已提交
411 412
TFileReader* tfileReaderOpen(char* path, uint64_t suid, int32_t version, const char* colName) {
  char fullname[256] = {0};
dengyihao's avatar
dengyihao 已提交
413 414
  tfileGenFileFullName(fullname, path, suid, colName, version);

dengyihao's avatar
dengyihao 已提交
415
  WriterCtx* wc = writerCtxCreate(TFile, fullname, true, 1024 * 1024 * 1024);
dengyihao's avatar
dengyihao 已提交
416
  indexInfo("open read file name:%s, file size: %d", wc->file.buf, wc->file.size);
dengyihao's avatar
dengyihao 已提交
417 418 419
  if (wc == NULL) {
    return NULL;
  }
dengyihao's avatar
dengyihao 已提交
420

dengyihao's avatar
dengyihao 已提交
421
  TFileReader* reader = tfileReaderCreate(wc);
dengyihao's avatar
dengyihao 已提交
422 423
  return reader;
}
dengyihao's avatar
dengyihao 已提交
424
TFileWriter* tfileWriterCreate(WriterCtx* ctx, TFileHeader* header) {
wafwerar's avatar
wafwerar 已提交
425
  TFileWriter* tw = taosMemoryCalloc(1, sizeof(TFileWriter));
dengyihao's avatar
dengyihao 已提交
426
  if (tw == NULL) {
dengyihao's avatar
dengyihao 已提交
427
    indexError("index: %" PRIu64 " failed to alloc TFilerWriter", header->suid);
dengyihao's avatar
dengyihao 已提交
428 429
    return NULL;
  }
dengyihao's avatar
dengyihao 已提交
430 431
  tw->ctx = ctx;
  tw->header = *header;
dengyihao's avatar
dengyihao 已提交
432
  tfileWriteHeader(tw);
dengyihao's avatar
dengyihao 已提交
433 434
  return tw;
}
dengyihao's avatar
dengyihao 已提交
435

dengyihao's avatar
dengyihao 已提交
436
int tfileWriterPut(TFileWriter* tw, void* data, bool order) {
dengyihao's avatar
dengyihao 已提交
437
  // sort by coltype and write to tindex
dengyihao's avatar
dengyihao 已提交
438 439
  if (order == false) {
    __compar_fn_t fn;
dengyihao's avatar
dengyihao 已提交
440 441

    int8_t colType = tw->header.colType;
dengyihao's avatar
add UT  
dengyihao 已提交
442
    colType = INDEX_TYPE_GET_TYPE(colType);
dengyihao's avatar
dengyihao 已提交
443 444 445 446 447 448
    if (colType == TSDB_DATA_TYPE_BINARY || colType == TSDB_DATA_TYPE_NCHAR) {
      fn = tfileStrCompare;
    } else {
      fn = getComparFunc(colType, 0);
    }
    taosArraySortPWithExt((SArray*)(data), tfileValueCompare, &fn);
dengyihao's avatar
dengyihao 已提交
449
  }
dengyihao's avatar
dengyihao 已提交
450

dengyihao's avatar
dengyihao 已提交
451
  int32_t bufLimit = 64 * 4096, offset = 0;
wafwerar's avatar
wafwerar 已提交
452
  // char*   buf = taosMemoryCalloc(1, sizeof(char) * bufLimit);
dengyihao's avatar
dengyihao 已提交
453
  // char*   p = buf;
dengyihao's avatar
dengyihao 已提交
454
  int32_t sz = taosArrayGetSize((SArray*)data);
dengyihao's avatar
dengyihao 已提交
455 456 457 458 459
  int32_t fstOffset = tw->offset;

  // ugly code, refactor later
  for (size_t i = 0; i < sz; i++) {
    TFileValue* v = taosArrayGetP((SArray*)data, i);
dengyihao's avatar
dengyihao 已提交
460 461
    taosArraySort(v->tableId, tfileUidCompare);
    taosArrayRemoveDuplicate(v->tableId, tfileUidCompare, NULL);
dengyihao's avatar
dengyihao 已提交
462
    int32_t tbsz = taosArrayGetSize(v->tableId);
dengyihao's avatar
dengyihao 已提交
463
    fstOffset += TF_TABLE_TATOAL_SIZE(tbsz);
dengyihao's avatar
dengyihao 已提交
464 465 466
  }
  tfileWriteFstOffset(tw, fstOffset);

dengyihao's avatar
dengyihao 已提交
467 468 469 470 471 472 473
  for (size_t i = 0; i < sz; i++) {
    TFileValue* v = taosArrayGetP((SArray*)data, i);

    int32_t tbsz = taosArrayGetSize(v->tableId);
    // check buf has enough space or not
    int32_t ttsz = TF_TABLE_TATOAL_SIZE(tbsz);

wafwerar's avatar
wafwerar 已提交
474
    char* buf = taosMemoryCalloc(1, ttsz * sizeof(char));
dengyihao's avatar
dengyihao 已提交
475
    char* p = buf;
dengyihao's avatar
dengyihao 已提交
476
    tfileSerialTableIdsToBuf(p, v->tableId);
dengyihao's avatar
dengyihao 已提交
477
    tw->ctx->write(tw->ctx, buf, ttsz);
dengyihao's avatar
dengyihao 已提交
478 479
    v->offset = tw->offset;
    tw->offset += ttsz;
wafwerar's avatar
wafwerar 已提交
480
    taosMemoryFree(buf);
dengyihao's avatar
dengyihao 已提交
481
  }
dengyihao's avatar
dengyihao 已提交
482

dengyihao's avatar
dengyihao 已提交
483 484
  tw->fb = fstBuilderCreate(tw->ctx, 0);
  if (tw->fb == NULL) {
dengyihao's avatar
dengyihao 已提交
485
    tfileWriterClose(tw);
dengyihao's avatar
dengyihao 已提交
486 487
    return -1;
  }
dengyihao's avatar
dengyihao 已提交
488 489

  // write data
dengyihao's avatar
dengyihao 已提交
490 491 492
  for (size_t i = 0; i < sz; i++) {
    // TODO, fst batch write later
    TFileValue* v = taosArrayGetP((SArray*)data, i);
dengyihao's avatar
dengyihao 已提交
493 494 495 496
    if (tfileWriteData(tw, v) != 0) {
      indexError("failed to write data: %s, offset: %d len: %d", v->colVal, v->offset,
                 (int)taosArrayGetSize(v->tableId));
    } else {
dengyihao's avatar
dengyihao 已提交
497 498
      // indexInfo("success to write data: %s, offset: %d len: %d", v->colVal, v->offset,
      //          (int)taosArrayGetSize(v->tableId));
dengyihao's avatar
dengyihao 已提交
499 500

      // indexInfo("tfile write data size: %d", tw->ctx->size(tw->ctx));
dengyihao's avatar
dengyihao 已提交
501 502
    }
  }
dengyihao's avatar
dengyihao 已提交
503 504 505
  fstBuilderFinish(tw->fb);
  fstBuilderDestroy(tw->fb);
  tw->fb = NULL;
dengyihao's avatar
dengyihao 已提交
506 507

  tfileWriteFooter(tw);
dengyihao's avatar
dengyihao 已提交
508 509
  return 0;
}
dengyihao's avatar
dengyihao 已提交
510
void tfileWriterClose(TFileWriter* tw) {
dengyihao's avatar
dengyihao 已提交
511 512 513
  if (tw == NULL) {
    return;
  }
dengyihao's avatar
dengyihao 已提交
514
  writerCtxDestroy(tw->ctx, false);
wafwerar's avatar
wafwerar 已提交
515
  taosMemoryFree(tw);
dengyihao's avatar
dengyihao 已提交
516
}
dengyihao's avatar
dengyihao 已提交
517
void tfileWriterDestroy(TFileWriter* tw) {
dengyihao's avatar
dengyihao 已提交
518 519 520
  if (tw == NULL) {
    return;
  }
dengyihao's avatar
dengyihao 已提交
521
  writerCtxDestroy(tw->ctx, false);
wafwerar's avatar
wafwerar 已提交
522
  taosMemoryFree(tw);
dengyihao's avatar
dengyihao 已提交
523
}
dengyihao's avatar
dengyihao 已提交
524

dengyihao's avatar
dengyihao 已提交
525
IndexTFile* indexTFileCreate(const char* path) {
dengyihao's avatar
dengyihao 已提交
526
  TFileCache* cache = tfileCacheCreate(path);
dengyihao's avatar
dengyihao 已提交
527 528 529
  if (cache == NULL) {
    return NULL;
  }
dengyihao's avatar
dengyihao 已提交
530

wafwerar's avatar
wafwerar 已提交
531
  IndexTFile* tfile = taosMemoryCalloc(1, sizeof(IndexTFile));
dengyihao's avatar
dengyihao 已提交
532 533 534 535
  if (tfile == NULL) {
    tfileCacheDestroy(cache);
    return NULL;
  }
536

dengyihao's avatar
dengyihao 已提交
537
  tfile->cache = cache;
dengyihao's avatar
dengyihao 已提交
538 539
  return tfile;
}
540
void indexTFileDestroy(IndexTFile* tfile) {
dengyihao's avatar
dengyihao 已提交
541 542 543
  if (tfile == NULL) {
    return;
  }
dengyihao's avatar
dengyihao 已提交
544
  tfileCacheDestroy(tfile->cache);
wafwerar's avatar
wafwerar 已提交
545
  taosMemoryFree(tfile);
dengyihao's avatar
dengyihao 已提交
546
}
dengyihao's avatar
dengyihao 已提交
547

dengyihao's avatar
dengyihao 已提交
548
int indexTFileSearch(void* tfile, SIndexTermQuery* query, SIdxTempResult* result) {
dengyihao's avatar
dengyihao 已提交
549
  int ret = -1;
dengyihao's avatar
dengyihao 已提交
550 551 552
  if (tfile == NULL) {
    return ret;
  }
dengyihao's avatar
dengyihao 已提交
553

dengyihao's avatar
add UT  
dengyihao 已提交
554
  int64_t     st = taosGetTimestampUs();
dengyihao's avatar
dengyihao 已提交
555
  IndexTFile* pTfile = tfile;
556

dengyihao's avatar
dengyihao 已提交
557 558
  SIndexTerm* term = query->term;
  ICacheKey key = {.suid = term->suid, .colType = term->colType, .colName = term->colName, .nColName = term->nColName};
559
  TFileReader* reader = tfileCacheGet(pTfile->cache, &key);
dengyihao's avatar
dengyihao 已提交
560 561 562
  if (reader == NULL) {
    return 0;
  }
dengyihao's avatar
add UT  
dengyihao 已提交
563 564
  int64_t cost = taosGetTimestampUs() - st;
  indexInfo("index tfile stage 1 cost: %" PRId64 "", cost);
dengyihao's avatar
dengyihao 已提交
565

dengyihao's avatar
dengyihao 已提交
566
  return tfileReaderSearch(reader, query, result);
dengyihao's avatar
dengyihao 已提交
567
}
dengyihao's avatar
dengyihao 已提交
568
int indexTFilePut(void* tfile, SIndexTerm* term, uint64_t uid) {
569 570
  // TFileWriterOpt wOpt = {.suid = term->suid, .colType = term->colType, .colName = term->colName, .nColName =
  // term->nColName, .version = 1};
dengyihao's avatar
dengyihao 已提交
571

572 573
  return 0;
}
dengyihao's avatar
dengyihao 已提交
574 575 576 577 578 579 580 581 582
static bool tfileIteratorNext(Iterate* iiter) {
  IterateValue* iv = &iiter->val;
  iterateValueDestroy(iv, false);

  char*    colVal = NULL;
  uint64_t offset = 0;

  TFileFstIter*          tIter = iiter->iter;
  StreamWithStateResult* rt = streamWithStateNextWith(tIter->st, NULL);
dengyihao's avatar
dengyihao 已提交
583 584 585
  if (rt == NULL) {
    return false;
  }
dengyihao's avatar
dengyihao 已提交
586 587 588

  int32_t sz = 0;
  char*   ch = (char*)fstSliceData(&rt->data, &sz);
wafwerar's avatar
wafwerar 已提交
589
  colVal = taosMemoryCalloc(1, sz + 1);
dengyihao's avatar
dengyihao 已提交
590 591 592 593 594
  memcpy(colVal, ch, sz);

  offset = (uint64_t)(rt->out.out);
  swsResultDestroy(rt);
  // set up iterate value
dengyihao's avatar
dengyihao 已提交
595 596 597
  if (tfileReaderLoadTableIds(tIter->rdr, offset, iv->val) != 0) {
    return false;
  }
dengyihao's avatar
dengyihao 已提交
598

dengyihao's avatar
dengyihao 已提交
599
  iv->ver = 0;
600
  iv->type = ADD_VALUE;  // value in tfile always ADD_VALUE
dengyihao's avatar
dengyihao 已提交
601
  iv->colVal = colVal;
dengyihao's avatar
dengyihao 已提交
602
  return true;
dengyihao's avatar
dengyihao 已提交
603 604 605
  // std::string key(ch, sz);
}

606
static IterateValue* tifileIterateGetValue(Iterate* iter) { return &iter->val; }
dengyihao's avatar
dengyihao 已提交
607 608

static TFileFstIter* tfileFstIteratorCreate(TFileReader* reader) {
wafwerar's avatar
wafwerar 已提交
609
  TFileFstIter* tIter = taosMemoryCalloc(1, sizeof(TFileFstIter));
dengyihao's avatar
dengyihao 已提交
610 611 612
  if (tIter == NULL) {
    return NULL;
  }
dengyihao's avatar
dengyihao 已提交
613

dengyihao's avatar
dengyihao 已提交
614 615 616 617 618 619 620 621
  tIter->ctx = automCtxCreate(NULL, AUTOMATION_ALWAYS);
  tIter->fb = fstSearch(reader->fst, tIter->ctx);
  tIter->st = streamBuilderIntoStream(tIter->fb);
  tIter->rdr = reader;
  return tIter;
}

Iterate* tfileIteratorCreate(TFileReader* reader) {
dengyihao's avatar
dengyihao 已提交
622 623 624
  if (reader == NULL) {
    return NULL;
  }
dengyihao's avatar
dengyihao 已提交
625

wafwerar's avatar
wafwerar 已提交
626
  Iterate* iter = taosMemoryCalloc(1, sizeof(Iterate));
dengyihao's avatar
dengyihao 已提交
627
  iter->iter = tfileFstIteratorCreate(reader);
628
  if (iter->iter == NULL) {
wafwerar's avatar
wafwerar 已提交
629
    taosMemoryFree(iter);
630 631
    return NULL;
  }
dengyihao's avatar
dengyihao 已提交
632 633
  iter->next = tfileIteratorNext;
  iter->getValue = tifileIterateGetValue;
dengyihao's avatar
dengyihao 已提交
634
  iter->val.val = taosArrayInit(1, sizeof(uint64_t));
635
  iter->val.colVal = NULL;
dengyihao's avatar
dengyihao 已提交
636 637 638
  return iter;
}
void tfileIteratorDestroy(Iterate* iter) {
dengyihao's avatar
dengyihao 已提交
639 640 641
  if (iter == NULL) {
    return;
  }
dengyihao's avatar
dengyihao 已提交
642

dengyihao's avatar
dengyihao 已提交
643 644 645 646 647 648 649
  IterateValue* iv = &iter->val;
  iterateValueDestroy(iv, true);

  TFileFstIter* tIter = iter->iter;
  streamWithStateDestroy(tIter->st);
  fstStreamBuilderDestroy(tIter->fb);
  automCtxDestroy(tIter->ctx);
wafwerar's avatar
wafwerar 已提交
650
  taosMemoryFree(tIter);
dengyihao's avatar
dengyihao 已提交
651

wafwerar's avatar
wafwerar 已提交
652
  taosMemoryFree(iter);
dengyihao's avatar
dengyihao 已提交
653 654
}

dengyihao's avatar
dengyihao 已提交
655
TFileReader* tfileGetReaderByCol(IndexTFile* tf, uint64_t suid, char* colName) {
dengyihao's avatar
dengyihao 已提交
656 657 658
  if (tf == NULL) {
    return NULL;
  }
dengyihao's avatar
dengyihao 已提交
659
  ICacheKey key = {.suid = suid, .colType = TSDB_DATA_TYPE_BINARY, .colName = colName, .nColName = strlen(colName)};
dengyihao's avatar
dengyihao 已提交
660 661
  return tfileCacheGet(tf->cache, &key);
}
dengyihao's avatar
dengyihao 已提交
662

dengyihao's avatar
dengyihao 已提交
663 664 665 666 667
static int tfileUidCompare(const void* a, const void* b) {
  uint64_t l = *(uint64_t*)a;
  uint64_t r = *(uint64_t*)b;
  return l - r;
}
dengyihao's avatar
dengyihao 已提交
668 669
static int tfileStrCompare(const void* a, const void* b) {
  int ret = strcmp((char*)a, (char*)b);
dengyihao's avatar
dengyihao 已提交
670 671 672
  if (ret == 0) {
    return ret;
  }
dengyihao's avatar
dengyihao 已提交
673 674 675
  return ret < 0 ? -1 : 1;
}

dengyihao's avatar
dengyihao 已提交
676 677 678 679 680 681 682 683
static int tfileValueCompare(const void* a, const void* b, const void* param) {
  __compar_fn_t fn = *(__compar_fn_t*)param;

  TFileValue* av = (TFileValue*)a;
  TFileValue* bv = (TFileValue*)b;

  return fn(av->colVal, bv->colVal);
}
dengyihao's avatar
dengyihao 已提交
684 685

TFileValue* tfileValueCreate(char* val) {
wafwerar's avatar
wafwerar 已提交
686
  TFileValue* tf = taosMemoryCalloc(1, sizeof(TFileValue));
dengyihao's avatar
dengyihao 已提交
687 688 689
  if (tf == NULL) {
    return NULL;
  }
690
  tf->colVal = tstrdup(val);
dengyihao's avatar
dengyihao 已提交
691 692 693 694
  tf->tableId = taosArrayInit(32, sizeof(uint64_t));
  return tf;
}
int tfileValuePush(TFileValue* tf, uint64_t val) {
dengyihao's avatar
dengyihao 已提交
695 696 697
  if (tf == NULL) {
    return -1;
  }
dengyihao's avatar
dengyihao 已提交
698 699 700 701 702
  taosArrayPush(tf->tableId, &val);
  return 0;
}
void tfileValueDestroy(TFileValue* tf) {
  taosArrayDestroy(tf->tableId);
wafwerar's avatar
wafwerar 已提交
703 704
  taosMemoryFree(tf->colVal);
  taosMemoryFree(tf);
dengyihao's avatar
dengyihao 已提交
705
}
dengyihao's avatar
dengyihao 已提交
706 707 708 709 710
static void tfileSerialTableIdsToBuf(char* buf, SArray* ids) {
  int sz = taosArrayGetSize(ids);
  SERIALIZE_VAR_TO_BUF(buf, sz, int32_t);
  for (size_t i = 0; i < sz; i++) {
    uint64_t* v = taosArrayGet(ids, i);
dengyihao's avatar
dengyihao 已提交
711 712 713 714 715 716 717
    SERIALIZE_VAR_TO_BUF(buf, *v, uint64_t);
  }
}

static int tfileWriteFstOffset(TFileWriter* tw, int32_t offset) {
  int32_t fstOffset = offset + sizeof(tw->header.fstOffset);
  tw->header.fstOffset = fstOffset;
dengyihao's avatar
dengyihao 已提交
718

dengyihao's avatar
dengyihao 已提交
719 720 721
  if (sizeof(fstOffset) != tw->ctx->write(tw->ctx, (char*)&fstOffset, sizeof(fstOffset))) {
    return -1;
  }
dengyihao's avatar
dengyihao 已提交
722
  indexInfo("tfile write fst offset: %d", tw->ctx->size(tw->ctx));
dengyihao's avatar
dengyihao 已提交
723
  tw->offset += sizeof(fstOffset);
dengyihao's avatar
dengyihao 已提交
724 725 726
  return 0;
}
static int tfileWriteHeader(TFileWriter* writer) {
dengyihao's avatar
dengyihao 已提交
727
  char buf[TFILE_HEADER_NO_FST] = {0};
dengyihao's avatar
dengyihao 已提交
728 729 730 731

  TFileHeader* header = &writer->header;
  memcpy(buf, (char*)header, sizeof(buf));

dengyihao's avatar
dengyihao 已提交
732
  indexInfo("tfile pre write header size: %d", writer->ctx->size(writer->ctx));
dengyihao's avatar
dengyihao 已提交
733
  int nwrite = writer->ctx->write(writer->ctx, buf, sizeof(buf));
dengyihao's avatar
dengyihao 已提交
734 735 736
  if (sizeof(buf) != nwrite) {
    return -1;
  }
dengyihao's avatar
dengyihao 已提交
737 738

  indexInfo("tfile after write header size: %d", writer->ctx->size(writer->ctx));
dengyihao's avatar
dengyihao 已提交
739 740 741 742 743 744
  writer->offset = nwrite;
  return 0;
}
static int tfileWriteData(TFileWriter* write, TFileValue* tval) {
  TFileHeader* header = &write->header;
  uint8_t      colType = header->colType;
dengyihao's avatar
add UT  
dengyihao 已提交
745 746

  colType = INDEX_TYPE_GET_TYPE(colType);
dengyihao's avatar
dengyihao 已提交
747 748 749 750 751 752 753 754 755 756 757
  if (colType == TSDB_DATA_TYPE_BINARY || colType == TSDB_DATA_TYPE_NCHAR) {
    FstSlice key = fstSliceCreate((uint8_t*)(tval->colVal), (size_t)strlen(tval->colVal));
    if (fstBuilderInsert(write->fb, key, tval->offset)) {
      fstSliceDestroy(&key);
      return 0;
    }
    fstSliceDestroy(&key);
    return -1;
  } else {
    // handle other type later
  }
dengyihao's avatar
dengyihao 已提交
758
  return 0;
dengyihao's avatar
dengyihao 已提交
759
}
dengyihao's avatar
dengyihao 已提交
760 761 762 763 764
static int tfileWriteFooter(TFileWriter* write) {
  char  buf[sizeof(tfileMagicNumber) + 1] = {0};
  void* pBuf = (void*)buf;
  taosEncodeFixedU64((void**)(void*)&pBuf, tfileMagicNumber);
  int nwrite = write->ctx->write(write->ctx, buf, strlen(buf));
dengyihao's avatar
dengyihao 已提交
765 766

  indexInfo("tfile write footer size: %d", write->ctx->size(write->ctx));
dengyihao's avatar
dengyihao 已提交
767 768 769
  assert(nwrite == sizeof(tfileMagicNumber));
  return nwrite;
}
dengyihao's avatar
dengyihao 已提交
770
static int tfileReaderLoadHeader(TFileReader* reader) {
dengyihao's avatar
dengyihao 已提交
771
  // TODO simple tfile header later
dengyihao's avatar
dengyihao 已提交
772
  char buf[TFILE_HEADER_SIZE] = {0};
dengyihao's avatar
dengyihao 已提交
773

dengyihao's avatar
dengyihao 已提交
774
  int64_t nread = reader->ctx->readFrom(reader->ctx, buf, sizeof(buf), 0);
dengyihao's avatar
dengyihao 已提交
775
  if (nread == -1) {
dengyihao's avatar
dengyihao 已提交
776 777
    indexError("actual Read: %d, to read: %d, errno: %d, filename: %s", (int)(nread), (int)sizeof(buf), errno,
               reader->ctx->file.buf);
dengyihao's avatar
dengyihao 已提交
778
  } else {
dengyihao's avatar
dengyihao 已提交
779
    indexInfo("actual Read: %d, to read: %d, filename: %s", (int)(nread), (int)sizeof(buf), reader->ctx->file.buf);
dengyihao's avatar
dengyihao 已提交
780 781
  }
  // assert(nread == sizeof(buf));
dengyihao's avatar
dengyihao 已提交
782
  memcpy(&reader->header, buf, sizeof(buf));
dengyihao's avatar
dengyihao 已提交
783

dengyihao's avatar
dengyihao 已提交
784 785
  return 0;
}
dengyihao's avatar
dengyihao 已提交
786
static int tfileReaderLoadFst(TFileReader* reader) {
dengyihao's avatar
dengyihao 已提交
787 788
  WriterCtx* ctx = reader->ctx;
  int        size = ctx->size(ctx);
dengyihao's avatar
dengyihao 已提交
789

dengyihao's avatar
dengyihao 已提交
790 791
  // current load fst into memory, refactor it later
  int   fstSize = size - reader->header.fstOffset - sizeof(tfileMagicNumber);
wafwerar's avatar
wafwerar 已提交
792
  char* buf = taosMemoryCalloc(1, fstSize);
dengyihao's avatar
dengyihao 已提交
793 794 795
  if (buf == NULL) {
    return -1;
  }
dengyihao's avatar
dengyihao 已提交
796

dengyihao's avatar
dengyihao 已提交
797
  int64_t ts = taosGetTimestampUs();
dengyihao's avatar
dengyihao 已提交
798
  int32_t nread = ctx->readFrom(ctx, buf, fstSize, reader->header.fstOffset);
dengyihao's avatar
dengyihao 已提交
799
  int64_t cost = taosGetTimestampUs() - ts;
dengyihao's avatar
dengyihao 已提交
800
  indexInfo("nread = %d, and fst offset=%d, fst size: %d, filename: %s, file size: %d, time cost: %" PRId64 "us", nread,
dengyihao's avatar
dengyihao 已提交
801
            reader->header.fstOffset, fstSize, ctx->file.buf, ctx->file.size, cost);
dengyihao's avatar
dengyihao 已提交
802
  // we assuse fst size less than FST_MAX_SIZE
dengyihao's avatar
dengyihao 已提交
803
  assert(nread > 0 && nread <= fstSize);
dengyihao's avatar
dengyihao 已提交
804 805 806

  FstSlice st = fstSliceCreate((uint8_t*)buf, nread);
  reader->fst = fstCreate(&st);
wafwerar's avatar
wafwerar 已提交
807
  taosMemoryFree(buf);
dengyihao's avatar
dengyihao 已提交
808 809
  fstSliceDestroy(&st);

dengyihao's avatar
dengyihao 已提交
810
  return reader->fst != NULL ? 0 : -1;
dengyihao's avatar
dengyihao 已提交
811
}
dengyihao's avatar
dengyihao 已提交
812
static int tfileReaderLoadTableIds(TFileReader* reader, int32_t offset, SArray* result) {
dengyihao's avatar
dengyihao 已提交
813
  // TODO(yihao): opt later
dengyihao's avatar
dengyihao 已提交
814
  WriterCtx* ctx = reader->ctx;
dengyihao's avatar
dengyihao 已提交
815 816 817 818 819 820 821 822 823 824 825 826 827 828 829 830
  char       block[1024] = {0};
  int32_t    nread = ctx->readFrom(ctx, block, sizeof(block), offset);
  assert(nread >= sizeof(uint32_t));

  char*   p = block;
  int32_t nid = *(int32_t*)p;
  p += sizeof(nid);

  while (nid > 0) {
    int32_t left = block + sizeof(block) - p;
    if (left >= sizeof(uint64_t)) {
      taosArrayPush(result, (uint64_t*)p);
      p += sizeof(uint64_t);
    } else {
      char buf[sizeof(uint64_t)] = {0};
      memcpy(buf, p, left);
dengyihao's avatar
dengyihao 已提交
831

dengyihao's avatar
dengyihao 已提交
832 833 834 835
      memset(block, 0, sizeof(block));
      offset += sizeof(block);
      nread = ctx->readFrom(ctx, block, sizeof(block), offset);
      memcpy(buf + left, block, sizeof(uint64_t) - left);
dengyihao's avatar
dengyihao 已提交
836

dengyihao's avatar
dengyihao 已提交
837 838 839 840
      taosArrayPush(result, (uint64_t*)buf);
      p = block + sizeof(uint64_t) - left;
    }
    nid -= 1;
dengyihao's avatar
dengyihao 已提交
841
  }
dengyihao's avatar
dengyihao 已提交
842 843
  return 0;
}
dengyihao's avatar
dengyihao 已提交
844 845 846 847 848 849 850 851 852 853 854 855 856 857 858 859 860 861 862
static int tfileReaderVerify(TFileReader* reader) {
  // just validate header and Footer, file corrupted also shuild be verified later
  WriterCtx* ctx = reader->ctx;

  uint64_t tMagicNumber = 0;

  char buf[sizeof(tMagicNumber) + 1] = {0};
  int  size = ctx->size(ctx);

  if (size < sizeof(tMagicNumber) || size <= sizeof(reader->header)) {
    return -1;
  } else if (ctx->readFrom(ctx, buf, sizeof(tMagicNumber), size - sizeof(tMagicNumber)) != sizeof(tMagicNumber)) {
    return -1;
  }

  taosDecodeFixedU64(buf, &tMagicNumber);
  return tMagicNumber == tfileMagicNumber ? 0 : -1;
}

dengyihao's avatar
dengyihao 已提交
863
void tfileReaderRef(TFileReader* reader) {
dengyihao's avatar
dengyihao 已提交
864 865 866
  if (reader == NULL) {
    return;
  }
dengyihao's avatar
dengyihao 已提交
867 868 869 870
  int ref = T_REF_INC(reader);
  UNUSED(ref);
}

dengyihao's avatar
dengyihao 已提交
871
void tfileReaderUnRef(TFileReader* reader) {
dengyihao's avatar
dengyihao 已提交
872 873 874
  if (reader == NULL) {
    return;
  }
dengyihao's avatar
dengyihao 已提交
875
  int ref = T_REF_DEC(reader);
876
  if (ref == 0) {
dengyihao's avatar
dengyihao 已提交
877
    // do nothing
878 879
    tfileReaderDestroy(reader);
  }
dengyihao's avatar
dengyihao 已提交
880
}
dengyihao's avatar
dengyihao 已提交
881

dengyihao's avatar
dengyihao 已提交
882
static SArray* tfileGetFileList(const char* path) {
dengyihao's avatar
dengyihao 已提交
883 884 885
  char     buf[128] = {0};
  uint64_t suid;
  uint32_t version;
dengyihao's avatar
dengyihao 已提交
886
  SArray*  files = taosArrayInit(4, sizeof(void*));
dengyihao's avatar
dengyihao 已提交
887

wafwerar's avatar
wafwerar 已提交
888 889
  TdDirPtr pDir = taosOpenDir(path);
  if (NULL == pDir) {
dengyihao's avatar
dengyihao 已提交
890 891
    return NULL;
  }
wafwerar's avatar
wafwerar 已提交
892 893 894
  TdDirEntryPtr pDirEntry;
  while ((pDirEntry = taosReadDir(pDir)) != NULL) {
    char* file = taosGetDirEntryName(pDirEntry);
dengyihao's avatar
dengyihao 已提交
895 896 897
    if (0 != tfileParseFileName(file, &suid, buf, &version)) {
      continue;
    }
dengyihao's avatar
dengyihao 已提交
898 899

    size_t len = strlen(path) + 1 + strlen(file) + 1;
wafwerar's avatar
wafwerar 已提交
900
    char*  buf = taosMemoryCalloc(1, len);
dengyihao's avatar
dengyihao 已提交
901
    sprintf(buf, "%s/%s", path, file);
dengyihao's avatar
dengyihao 已提交
902
    taosArrayPush(files, &buf);
dengyihao's avatar
dengyihao 已提交
903
  }
wafwerar's avatar
wafwerar 已提交
904
  taosCloseDir(&pDir);
dengyihao's avatar
dengyihao 已提交
905 906 907 908 909

  taosArraySort(files, tfileCompare);
  tfileRmExpireFile(files);

  return files;
dengyihao's avatar
dengyihao 已提交
910
}
dengyihao's avatar
dengyihao 已提交
911 912 913 914
static int tfileRmExpireFile(SArray* result) {
  // TODO(yihao): remove expire tindex after restart
  return 0;
}
dengyihao's avatar
dengyihao 已提交
915 916
static void tfileDestroyFileName(void* elem) {
  char* p = *(char**)elem;
wafwerar's avatar
wafwerar 已提交
917
  taosMemoryFree(p);
dengyihao's avatar
dengyihao 已提交
918 919
}
static int tfileCompare(const void* a, const void* b) {
dengyihao's avatar
dengyihao 已提交
920 921 922
  const char* as = *(char**)a;
  const char* bs = *(char**)b;
  return strcmp(as, bs);
dengyihao's avatar
dengyihao 已提交
923
}
dengyihao's avatar
dengyihao 已提交
924 925 926

static int tfileParseFileName(const char* filename, uint64_t* suid, char* col, int* version) {
  if (3 == sscanf(filename, "%" PRIu64 "-%[^-]-%d.tindex", suid, col, version)) {
dengyihao's avatar
dengyihao 已提交
927 928 929 930 931
    // read suid & colid & version  success
    return 0;
  }
  return -1;
}
dengyihao's avatar
dengyihao 已提交
932 933 934 935 936 937 938 939 940 941
// tfile name suid-colId-version.tindex
static void tfileGenFileName(char* filename, uint64_t suid, const char* col, int version) {
  sprintf(filename, "%" PRIu64 "-%s-%d.tindex", suid, col, version);
  return;
}
static void tfileGenFileFullName(char* fullname, const char* path, uint64_t suid, const char* col, int32_t version) {
  char filename[128] = {0};
  tfileGenFileName(filename, suid, col, version);
  sprintf(fullname, "%s/%s", path, filename);
}