streamBackendRocksdb.c 81.8 KB
Newer Older
dengyihao's avatar
dengyihao 已提交
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15
/*
 * Copyright (c) 2019 TAOS Data, Inc. <jhtao@taosdata.com>
 *
 * 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 "streamBackendRocksdb.h"
dengyihao's avatar
dengyihao 已提交
17
#include "executor.h"
dengyihao's avatar
dengyihao 已提交
18
#include "query.h"
dengyihao's avatar
dengyihao 已提交
19
#include "streamInc.h"
dengyihao's avatar
dengyihao 已提交
20
#include "tcommon.h"
dengyihao's avatar
dengyihao 已提交
21
#include "tref.h"
dengyihao's avatar
dengyihao 已提交
22

23 24 25 26
typedef struct SCompactFilteFactory {
  void* status;
} SCompactFilteFactory;

Y
yihaoDeng 已提交
27 28 29
typedef struct {
  void* tableOpt;
} RocksdbCfParam;
dengyihao's avatar
dengyihao 已提交
30 31 32 33 34 35 36
typedef struct {
  rocksdb_t*                       db;
  rocksdb_column_family_handle_t** pHandle;
  rocksdb_writeoptions_t*          wOpt;
  rocksdb_readoptions_t*           rOpt;
  rocksdb_options_t**              cfOpt;
  rocksdb_options_t*               dbOpt;
Y
yihaoDeng 已提交
37 38
  RocksdbCfParam*                  param;
  void*                            pBackend;
dengyihao's avatar
dengyihao 已提交
39
  SListNode*                       pCompareNode;
Y
yihaoDeng 已提交
40
  rocksdb_comparator_t**           pCompares;
dengyihao's avatar
dengyihao 已提交
41 42
} RocksdbCfInst;

Y
yihaoDeng 已提交
43 44
uint32_t nextPow2(uint32_t x);
int32_t  streamStateOpenBackendCf(void* backend, char* name, char** cfs, int32_t nCf);
dengyihao's avatar
dengyihao 已提交
45 46

void destroyRocksdbCfInst(RocksdbCfInst* inst);
dengyihao's avatar
dengyihao 已提交
47

48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65
void          destroyCompactFilteFactory(void* arg);
void          destroyCompactFilte(void* arg);
const char*   compactFilteFactoryName(void* arg);
const char*   compactFilteName(void* arg);
unsigned char compactFilte(void* arg, int level, const char* key, size_t klen, const char* val, size_t vlen,
                           char** newval, size_t* newvlen, unsigned char* value_changed);
rocksdb_compactionfilter_t* compactFilteFactoryCreateFilter(void* arg, rocksdb_compactionfiltercontext_t* ctx);

const char* cfName[] = {"default", "state", "fill", "sess", "func", "parname", "partag"};

typedef int (*EncodeFunc)(void* key, char* buf);
typedef int (*DecodeFunc)(void* key, char* buf);
typedef int (*ToStringFunc)(void* key, char* buf);
typedef const char* (*CompareName)(void* statue);
typedef int (*BackendCmpFunc)(void* state, const char* aBuf, size_t aLen, const char* bBuf, size_t bLen);
typedef void (*DestroyFunc)(void* state);
typedef int32_t (*EncodeValueFunc)(void* value, int32_t vlen, int64_t ttl, char** dest);
typedef int32_t (*DecodeValueFunc)(void* value, int32_t vlen, int64_t* ttl, char** dest);
Y
yihaoDeng 已提交
66 67 68 69 70 71 72 73 74 75 76 77
typedef struct {
  const char*     key;
  int32_t         len;
  int             idx;
  BackendCmpFunc  cmpFunc;
  EncodeFunc      enFunc;
  DecodeFunc      deFunc;
  ToStringFunc    toStrFunc;
  CompareName     cmpName;
  DestroyFunc     detroyFunc;
  EncodeValueFunc enValueFunc;
  DecodeValueFunc deValueFunc;
78

Y
yihaoDeng 已提交
79
} SCfInit;
80

Y
yihaoDeng 已提交
81
#define GEN_COLUMN_FAMILY_NAME(name, idstr, SUFFIX) sprintf(name, "%s_%s", idstr, (SUFFIX));
82 83 84 85 86 87 88 89
const char* compareDefaultName(void* name);
const char* compareStateName(void* name);
const char* compareWinKeyName(void* name);
const char* compareSessionKeyName(void* name);
const char* compareFuncKeyName(void* name);
const char* compareParKeyName(void* name);
const char* comparePartagKeyName(void* name);

Y
yihaoDeng 已提交
90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145
int defaultKeyComp(void* state, const char* aBuf, size_t aLen, const char* bBuf, size_t bLen);
int defaultKeyEncode(void* k, char* buf);
int defaultKeyDecode(void* k, char* buf);
int defaultKeyToString(void* k, char* buf);

int stateKeyDBComp(void* state, const char* aBuf, size_t aLen, const char* bBuf, size_t bLen);
int stateKeyEncode(void* k, char* buf);
int stateKeyDecode(void* k, char* buf);
int stateKeyToString(void* k, char* buf);

int stateSessionKeyDBComp(void* state, const char* aBuf, size_t aLen, const char* bBuf, size_t bLen);
int stateSessionKeyEncode(void* ses, char* buf);
int stateSessionKeyDecode(void* ses, char* buf);
int stateSessionKeyToString(void* k, char* buf);

int winKeyDBComp(void* state, const char* aBuf, size_t aLen, const char* bBuf, size_t bLen);
int winKeyEncode(void* k, char* buf);
int winKeyDecode(void* k, char* buf);
int winKeyToString(void* k, char* buf);

int tupleKeyDBComp(void* state, const char* aBuf, size_t aLen, const char* bBuf, size_t bLen);
int tupleKeyEncode(void* k, char* buf);
int tupleKeyDecode(void* k, char* buf);
int tupleKeyToString(void* k, char* buf);

int parKeyDBComp(void* state, const char* aBuf, size_t aLen, const char* bBuf, size_t bLen);
int parKeyEncode(void* k, char* buf);
int parKeyDecode(void* k, char* buf);
int parKeyToString(void* k, char* buf);

int     stremaValueEncode(void* k, char* buf);
int     streamValueDecode(void* k, char* buf);
int32_t streamValueToString(void* k, char* buf);
int32_t streaValueIsStale(void* k, int64_t ts);
void    destroyFunc(void* arg);

int32_t encodeValueFunc(void* value, int32_t vlen, int64_t ttl, char** dest);
int32_t decodeValueFunc(void* value, int32_t vlen, int64_t* ttl, char** dest);

SCfInit ginitDict[] = {
    {"default", 7, 0, defaultKeyComp, defaultKeyEncode, defaultKeyDecode, defaultKeyToString, compareDefaultName,
     destroyFunc, encodeValueFunc, decodeValueFunc},
    {"state", 5, 1, stateKeyDBComp, stateKeyEncode, stateKeyDecode, stateKeyToString, compareStateName, destroyFunc,
     encodeValueFunc, decodeValueFunc},
    {"fill", 4, 2, winKeyDBComp, winKeyEncode, winKeyDecode, winKeyToString, compareWinKeyName, destroyFunc,
     encodeValueFunc, decodeValueFunc},
    {"sess", 4, 3, stateSessionKeyDBComp, stateSessionKeyEncode, stateSessionKeyDecode, stateSessionKeyToString,
     compareSessionKeyName, destroyFunc, encodeValueFunc, decodeValueFunc},
    {"func", 4, 4, tupleKeyDBComp, tupleKeyEncode, tupleKeyDecode, tupleKeyToString, compareFuncKeyName, destroyFunc,
     encodeValueFunc, decodeValueFunc},
    {"parname", 7, 5, parKeyDBComp, parKeyEncode, parKeyDecode, parKeyToString, compareParKeyName, destroyFunc,
     encodeValueFunc, decodeValueFunc},
    {"partag", 6, 6, parKeyDBComp, parKeyEncode, parKeyDecode, parKeyToString, comparePartagKeyName, destroyFunc,
     encodeValueFunc, decodeValueFunc},
};

146
void* streamBackendInit(const char* path) {
dengyihao's avatar
dengyihao 已提交
147 148
  uint32_t dbMemLimit = nextPow2(tsMaxStreamBackendCache) << 20;

dengyihao's avatar
dengyihao 已提交
149
  qDebug("start to init stream backend at %s", path);
Y
yihaoDeng 已提交
150
  SBackendWrapper* pHandle = taosMemoryCalloc(1, sizeof(SBackendWrapper));
151 152
  pHandle->list = tdListNew(sizeof(SCfComparator));
  taosThreadMutexInit(&pHandle->mutex, NULL);
dengyihao's avatar
dengyihao 已提交
153 154
  taosThreadMutexInit(&pHandle->cfMutex, NULL);
  pHandle->cfInst = taosHashInit(64, taosGetDefaultHashFunction(TSDB_DATA_TYPE_BINARY), false, HASH_NO_LOCK);
155 156 157

  rocksdb_env_t* env = rocksdb_create_default_env();  // rocksdb_envoptions_create();

dengyihao's avatar
dengyihao 已提交
158 159 160 161 162
  int32_t nBGThread = tsNumOfSnodeStreamThreads <= 2 ? 1 : tsNumOfSnodeStreamThreads / 2;
  rocksdb_env_set_low_priority_background_threads(env, nBGThread);
  rocksdb_env_set_high_priority_background_threads(env, nBGThread);

  rocksdb_cache_t* cache = rocksdb_cache_create_lru(dbMemLimit / 2);
163 164 165 166 167

  rocksdb_options_t* opts = rocksdb_options_create();
  rocksdb_options_set_env(opts, env);
  rocksdb_options_set_create_if_missing(opts, 1);
  rocksdb_options_set_create_missing_column_families(opts, 1);
dengyihao's avatar
dengyihao 已提交
168
  rocksdb_options_set_max_total_wal_size(opts, dbMemLimit);
169
  rocksdb_options_set_recycle_log_file_num(opts, 6);
dengyihao's avatar
dengyihao 已提交
170
  rocksdb_options_set_max_write_buffer_number(opts, 3);
dengyihao's avatar
dengyihao 已提交
171
  rocksdb_options_set_info_log_level(opts, 0);
dengyihao's avatar
dengyihao 已提交
172 173
  rocksdb_options_set_db_write_buffer_size(opts, dbMemLimit);
  rocksdb_options_set_write_buffer_size(opts, dbMemLimit / 2);
174 175 176 177 178 179 180 181

  pHandle->env = env;
  pHandle->dbOpt = opts;
  pHandle->cache = cache;
  pHandle->filterFactory = rocksdb_compactionfilterfactory_create(
      NULL, destroyCompactFilteFactory, compactFilteFactoryCreateFilter, compactFilteFactoryName);
  rocksdb_options_set_compaction_filter_factory(pHandle->dbOpt, pHandle->filterFactory);

dengyihao's avatar
dengyihao 已提交
182 183 184 185
  char*  err = NULL;
  size_t nCf = 0;

  char** cfs = rocksdb_list_column_families(opts, path, &nCf, &err);
dengyihao's avatar
dengyihao 已提交
186
  if (nCf == 0 || nCf == 1 || err != NULL) {
187
    taosMemoryFreeClear(err);
dengyihao's avatar
dengyihao 已提交
188 189 190 191
    pHandle->db = rocksdb_open(opts, path, &err);
    if (err != NULL) {
      qError("failed to open rocksdb, path:%s, reason:%s", path, err);
      taosMemoryFreeClear(err);
dengyihao's avatar
dengyihao 已提交
192
      goto _EXIT;
dengyihao's avatar
dengyihao 已提交
193 194
    }
  } else {
dengyihao's avatar
dengyihao 已提交
195 196 197
    /*
      list all cf and get prefix
    */
Y
yihaoDeng 已提交
198
    streamStateOpenBackendCf(pHandle, (char*)path, cfs, nCf);
199
  }
200 201 202
  if (cfs != NULL) {
    rocksdb_list_column_families_destroy(cfs, nCf);
  }
dengyihao's avatar
dengyihao 已提交
203
  qDebug("succ to init stream backend at %s, backend:%p", path, pHandle);
204 205 206 207 208 209 210

  return (void*)pHandle;
_EXIT:
  rocksdb_options_destroy(opts);
  rocksdb_cache_destroy(cache);
  rocksdb_env_destroy(env);
  taosThreadMutexDestroy(&pHandle->mutex);
dengyihao's avatar
dengyihao 已提交
211 212
  taosThreadMutexDestroy(&pHandle->cfMutex);
  taosHashCleanup(pHandle->cfInst);
213 214
  rocksdb_compactionfilterfactory_destroy(pHandle->filterFactory);
  tdListFree(pHandle->list);
dengyihao's avatar
dengyihao 已提交
215
  taosMemoryFree(pHandle);
dengyihao's avatar
dengyihao 已提交
216
  qDebug("failed to init stream backend at %s", path);
217 218 219
  return NULL;
}
void streamBackendCleanup(void* arg) {
Y
yihaoDeng 已提交
220 221
  SBackendWrapper* pHandle = (SBackendWrapper*)arg;
  RocksdbCfInst**  pIter = (RocksdbCfInst**)taosHashIterate(pHandle->cfInst, NULL);
dengyihao's avatar
dengyihao 已提交
222 223 224 225 226 227
  while (pIter != NULL) {
    RocksdbCfInst* inst = *pIter;
    destroyRocksdbCfInst(inst);
    taosHashIterate(pHandle->cfInst, pIter);
  }
  taosHashCleanup(pHandle->cfInst);
dengyihao's avatar
dengyihao 已提交
228

Y
yihaoDeng 已提交
229 230 231 232 233 234 235 236 237 238
  if (pHandle->db) {
    char*                   err = NULL;
    rocksdb_flushoptions_t* flushOpt = rocksdb_flushoptions_create();
    rocksdb_flush(pHandle->db, flushOpt, &err);
    if (err != NULL) {
      qError("failed to flush db before streamBackend clean up, reason:%s", err);
      taosMemoryFree(err);
    }
    rocksdb_flushoptions_destroy(flushOpt);
    rocksdb_close(pHandle->db);
dengyihao's avatar
dengyihao 已提交
239
  }
240 241 242 243 244 245 246 247 248 249
  rocksdb_options_destroy(pHandle->dbOpt);
  rocksdb_env_destroy(pHandle->env);
  rocksdb_cache_destroy(pHandle->cache);

  SListNode* head = tdListPopHead(pHandle->list);
  while (head != NULL) {
    streamStateDestroyCompar(head->data);
    taosMemoryFree(head);
    head = tdListPopHead(pHandle->list);
  }
dengyihao's avatar
dengyihao 已提交
250

251
  tdListFree(pHandle->list);
dengyihao's avatar
dengyihao 已提交
252 253
  taosThreadMutexDestroy(&pHandle->mutex);

dengyihao's avatar
dengyihao 已提交
254
  taosThreadMutexDestroy(&pHandle->cfMutex);
255 256

  taosMemoryFree(pHandle);
dengyihao's avatar
dengyihao 已提交
257
  qDebug("destroy stream backend backend:%p", pHandle);
258 259
  return;
}
Y
yihaoDeng 已提交
260
void streamBackendHandleCleanup(void* arg) {
Y
yihaoDeng 已提交
261 262
  SBackendCfWrapper* wrapper = arg;
  bool               remove = wrapper->remove;
Y
yihaoDeng 已提交
263 264 265 266 267 268 269 270
  qDebug("start to do-close backendwrapper %p, %s", wrapper, wrapper->idstr);
  if (wrapper->rocksdb == NULL) {
    return;
  }

  int cfLen = sizeof(ginitDict) / sizeof(ginitDict[0]);

  char* err = NULL;
Y
yihaoDeng 已提交
271
  if (remove) {
Y
yihaoDeng 已提交
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
    for (int i = 0; i < cfLen; i++) {
      if (wrapper->pHandle[i] != NULL)
        rocksdb_drop_column_family(wrapper->rocksdb, ((rocksdb_column_family_handle_t**)wrapper->pHandle)[i], &err);
      if (err != NULL) {
        // qError("failed to create cf:%s_%s, reason:%s", wrapper->idstr, ginitDict[i].key, err);
        taosMemoryFreeClear(err);
      }
    }
  } else {
    rocksdb_flushoptions_t* flushOpt = rocksdb_flushoptions_create();
    for (int i = 0; i < cfLen; i++) {
      if (wrapper->pHandle[i] != NULL) rocksdb_flush_cf(wrapper->rocksdb, flushOpt, wrapper->pHandle[i], &err);
      if (err != NULL) {
        qError("failed to create cf:%s_%s, reason:%s", wrapper->idstr, ginitDict[i].key, err);
        taosMemoryFreeClear(err);
      }
    }
    rocksdb_flushoptions_destroy(flushOpt);
  }

  for (int i = 0; i < cfLen; i++) {
    if (wrapper->pHandle[i] != NULL) {
      rocksdb_column_family_handle_destroy(wrapper->pHandle[i]);
    }
  }
  taosMemoryFreeClear(wrapper->pHandle);
  for (int i = 0; i < cfLen; i++) {
    rocksdb_options_destroy(wrapper->cfOpts[i]);
    rocksdb_block_based_options_destroy(((RocksdbCfParam*)wrapper->param)[i].tableOpt);
  }

Y
yihaoDeng 已提交
303
  if (remove) {
Y
yihaoDeng 已提交
304 305 306 307 308 309 310 311 312 313 314 315 316 317 318 319 320 321
    streamBackendDelCompare(wrapper->pBackend, wrapper->pComparNode);
  }
  rocksdb_writeoptions_destroy(wrapper->writeOpts);
  wrapper->writeOpts = NULL;

  rocksdb_readoptions_destroy(wrapper->readOpts);
  wrapper->readOpts = NULL;
  taosMemoryFreeClear(wrapper->cfOpts);
  taosMemoryFreeClear(wrapper->param);

  taosThreadRwlockDestroy(&wrapper->rwLock);
  wrapper->rocksdb = NULL;
  taosReleaseRef(streamBackendId, wrapper->backendId);

  qDebug("end to do-close backendwrapper %p, %s", wrapper, wrapper->idstr);
  taosMemoryFree(wrapper);
  return;
}
322
SListNode* streamBackendAddCompare(void* backend, void* arg) {
Y
yihaoDeng 已提交
323 324
  SBackendWrapper* pHandle = (SBackendWrapper*)backend;
  SListNode*       node = NULL;
325 326 327 328 329 330
  taosThreadMutexLock(&pHandle->mutex);
  node = tdListAdd(pHandle->list, arg);
  taosThreadMutexUnlock(&pHandle->mutex);
  return node;
}
void streamBackendDelCompare(void* backend, void* arg) {
Y
yihaoDeng 已提交
331 332
  SBackendWrapper* pHandle = (SBackendWrapper*)backend;
  SListNode*       node = NULL;
333 334 335 336 337 338 339 340 341
  taosThreadMutexLock(&pHandle->mutex);
  node = tdListPopNode(pHandle->list, arg);
  taosThreadMutexUnlock(&pHandle->mutex);
  if (node) {
    streamStateDestroyCompar(node->data);
    taosMemoryFree(node);
  }
}
void        streamStateDestroy_rocksdb(SStreamState* pState, bool remove) { streamStateCloseBackend(pState, remove); }
dengyihao's avatar
dengyihao 已提交
342 343 344 345 346 347 348 349
static bool streamStateIterSeekAndValid(rocksdb_iterator_t* iter, char* buf, size_t len);

// |key|-----value------|
// |key|ttl|len|userData|

static rocksdb_iterator_t* streamStateIterCreate(SStreamState* pState, const char* cfName,
                                                 rocksdb_snapshot_t** snapshot, rocksdb_readoptions_t** readOpt);

dengyihao's avatar
dengyihao 已提交
350
int defaultKeyComp(void* state, const char* aBuf, size_t aLen, const char* bBuf, size_t bLen) {
dengyihao's avatar
dengyihao 已提交
351 352 353 354 355 356 357 358 359 360 361
  int ret = memcmp(aBuf, bBuf, aLen);
  if (ret == 0) {
    if (aLen < bLen)
      return -1;
    else if (aLen > bLen)
      return 1;
    else
      return 0;
  } else {
    return ret;
  }
dengyihao's avatar
dengyihao 已提交
362
}
dengyihao's avatar
dengyihao 已提交
363 364 365
int streamStateValueIsStale(char* vv) {
  int64_t ts = 0;
  taosDecodeFixedI64(vv, &ts);
dengyihao's avatar
dengyihao 已提交
366
  return (ts != 0 && ts < taosGetTimestampMs()) ? 1 : 0;
dengyihao's avatar
dengyihao 已提交
367
}
dengyihao's avatar
dengyihao 已提交
368
int iterValueIsStale(rocksdb_iterator_t* iter) {
dengyihao's avatar
dengyihao 已提交
369 370 371
  size_t len;
  char*  v = (char*)rocksdb_iter_value(iter, &len);
  return streamStateValueIsStale(v);
dengyihao's avatar
dengyihao 已提交
372
}
dengyihao's avatar
dengyihao 已提交
373 374 375 376 377 378 379 380 381 382 383 384 385 386
int defaultKeyEncode(void* k, char* buf) {
  int len = strlen((char*)k);
  memcpy(buf, (char*)k, len);
  return len;
}
int defaultKeyDecode(void* k, char* buf) {
  int len = strlen(buf);
  memcpy(k, buf, len);
  return len;
}
int defaultKeyToString(void* k, char* buf) {
  // just to debug
  return sprintf(buf, "key: %s", (char*)k);
}
dengyihao's avatar
dengyihao 已提交
387 388 389 390 391 392 393 394 395 396 397 398 399 400 401 402 403 404 405 406 407 408 409 410 411 412 413 414 415 416 417 418
//
//  SStateKey
//  |--groupid--|---ts------|--opNum----|
//  |--uint64_t-|-uint64_t--|--int64_t--|
//
//
//
int stateKeyDBComp(void* state, const char* aBuf, size_t aLen, const char* bBuf, size_t bLen) {
  SStateKey key1, key2;
  memset(&key1, 0, sizeof(key1));
  memset(&key2, 0, sizeof(key2));

  char* p1 = (char*)aBuf;
  char* p2 = (char*)bBuf;

  p1 = taosDecodeFixedU64(p1, &key1.key.groupId);
  p2 = taosDecodeFixedU64(p2, &key2.key.groupId);

  p1 = taosDecodeFixedI64(p1, &key1.key.ts);
  p2 = taosDecodeFixedI64(p2, &key2.key.ts);

  taosDecodeFixedI64(p1, &key1.opNum);
  taosDecodeFixedI64(p2, &key2.opNum);

  return stateKeyCmpr(&key1, sizeof(key1), &key2, sizeof(key2));
}

int stateKeyEncode(void* k, char* buf) {
  SStateKey* key = k;
  int        len = 0;
  len += taosEncodeFixedU64((void**)&buf, key->key.groupId);
  len += taosEncodeFixedI64((void**)&buf, key->key.ts);
dengyihao's avatar
dengyihao 已提交
419
  len += taosEncodeFixedI64((void**)&buf, key->opNum);
dengyihao's avatar
dengyihao 已提交
420 421 422 423 424 425 426 427 428 429 430 431
  return len;
}
int stateKeyDecode(void* k, char* buf) {
  SStateKey* key = k;
  int        len = 0;
  char*      p = buf;
  p = taosDecodeFixedU64(p, &key->key.groupId);
  p = taosDecodeFixedI64(p, &key->key.ts);
  p = taosDecodeFixedI64(p, &key->opNum);
  return p - buf;
}

dengyihao's avatar
dengyihao 已提交
432 433 434
int stateKeyToString(void* k, char* buf) {
  SStateKey* key = k;
  int        n = 0;
dengyihao's avatar
dengyihao 已提交
435 436 437
  n += sprintf(buf + n, "[groupId:%" PRId64 ",", key->key.groupId);
  n += sprintf(buf + n, "ts:%" PRIi64 ",", key->key.ts);
  n += sprintf(buf + n, "opNum:%" PRIi64 "]", key->opNum);
dengyihao's avatar
dengyihao 已提交
438 439 440
  return n;
}

dengyihao's avatar
dengyihao 已提交
441 442 443 444 445 446 447 448 449 450 451 452 453 454 455 456 457 458 459 460 461 462 463 464 465 466 467 468 469 470 471 472 473 474 475 476 477 478 479 480 481 482 483 484 485 486 487 488 489 490
//
// SStateSessionKey
//  |-----------SSessionKey----------|
//  |-----STimeWindow-----|
//  |---skey--|---ekey----|--groupId-|--opNum--|
//  |---int64-|--int64_t--|--uint64--|--int64_t|
// |
//
int stateSessionKeyDBComp(void* state, const char* aBuf, size_t aLen, const char* bBuf, size_t bLen) {
  SStateSessionKey w1, w2;
  memset(&w1, 0, sizeof(w1));
  memset(&w2, 0, sizeof(w2));

  char* p1 = (char*)aBuf;
  char* p2 = (char*)bBuf;

  p1 = taosDecodeFixedI64(p1, &w1.key.win.skey);
  p2 = taosDecodeFixedI64(p2, &w2.key.win.skey);

  p1 = taosDecodeFixedI64(p1, &w1.key.win.ekey);
  p2 = taosDecodeFixedI64(p2, &w2.key.win.ekey);

  p1 = taosDecodeFixedU64(p1, &w1.key.groupId);
  p2 = taosDecodeFixedU64(p2, &w2.key.groupId);

  p1 = taosDecodeFixedI64(p1, &w1.opNum);
  p2 = taosDecodeFixedI64(p2, &w2.opNum);

  return stateSessionKeyCmpr(&w1, sizeof(w1), &w2, sizeof(w2));
}
int stateSessionKeyEncode(void* ses, char* buf) {
  SStateSessionKey* sess = ses;
  int               len = 0;
  len += taosEncodeFixedI64((void**)&buf, sess->key.win.skey);
  len += taosEncodeFixedI64((void**)&buf, sess->key.win.ekey);
  len += taosEncodeFixedU64((void**)&buf, sess->key.groupId);
  len += taosEncodeFixedI64((void**)&buf, sess->opNum);
  return len;
}
int stateSessionKeyDecode(void* ses, char* buf) {
  SStateSessionKey* sess = ses;
  int               len = 0;

  char* p = buf;
  p = taosDecodeFixedI64(p, &sess->key.win.skey);
  p = taosDecodeFixedI64(p, &sess->key.win.ekey);
  p = taosDecodeFixedU64(p, &sess->key.groupId);
  p = taosDecodeFixedI64(p, &sess->opNum);
  return p - buf;
}
dengyihao's avatar
dengyihao 已提交
491 492 493
int stateSessionKeyToString(void* k, char* buf) {
  SStateSessionKey* key = k;
  int               n = 0;
dengyihao's avatar
dengyihao 已提交
494 495 496 497
  n += sprintf(buf + n, "[skey:%" PRIi64 ",", key->key.win.skey);
  n += sprintf(buf + n, "ekey:%" PRIi64 ",", key->key.win.ekey);
  n += sprintf(buf + n, "groupId:%" PRIu64 ",", key->key.groupId);
  n += sprintf(buf + n, "opNum:%" PRIi64 "]", key->opNum);
dengyihao's avatar
dengyihao 已提交
498 499
  return n;
}
dengyihao's avatar
dengyihao 已提交
500 501 502 503 504 505 506 507 508 509 510 511 512 513 514 515 516 517 518 519 520 521 522 523 524 525 526 527 528 529 530 531 532 533 534 535 536 537 538

/**
 *  SWinKey
 *  |------groupId------|-----ts------|
 *  |------uint64-------|----int64----|
 */
int winKeyDBComp(void* state, const char* aBuf, size_t aLen, const char* bBuf, size_t bLen) {
  SWinKey w1, w2;
  memset(&w1, 0, sizeof(w1));
  memset(&w2, 0, sizeof(w2));

  char* p1 = (char*)aBuf;
  char* p2 = (char*)bBuf;

  p1 = taosDecodeFixedU64(p1, &w1.groupId);
  p2 = taosDecodeFixedU64(p2, &w2.groupId);

  p1 = taosDecodeFixedI64(p1, &w1.ts);
  p2 = taosDecodeFixedI64(p2, &w2.ts);

  return winKeyCmpr(&w1, sizeof(w1), &w2, sizeof(w2));
}

int winKeyEncode(void* k, char* buf) {
  SWinKey* key = k;
  int      len = 0;
  len += taosEncodeFixedU64((void**)&buf, key->groupId);
  len += taosEncodeFixedI64((void**)&buf, key->ts);
  return len;
}

int winKeyDecode(void* k, char* buf) {
  SWinKey* key = k;
  int      len = 0;
  char*    p = buf;
  p = taosDecodeFixedU64(p, &key->groupId);
  p = taosDecodeFixedI64(p, &key->ts);
  return len;
}
dengyihao's avatar
dengyihao 已提交
539 540 541 542

int winKeyToString(void* k, char* buf) {
  SWinKey* key = k;
  int      n = 0;
dengyihao's avatar
dengyihao 已提交
543 544
  n += sprintf(buf + n, "[groupId:%" PRIu64 ",", key->groupId);
  n += sprintf(buf + n, "ts:%" PRIi64 "]", key->ts);
dengyihao's avatar
dengyihao 已提交
545 546
  return n;
}
dengyihao's avatar
dengyihao 已提交
547 548 549 550 551 552 553 554 555 556 557 558 559 560 561 562 563 564 565 566 567 568 569 570 571 572 573 574 575 576 577 578 579 580 581 582 583 584 585 586 587 588
/*
 * STupleKey
 * |---groupId---|---ts---|---exprIdx---|
 * |---uint64--|---int64--|---int32-----|
 */
int tupleKeyDBComp(void* state, const char* aBuf, size_t aLen, const char* bBuf, size_t bLen) {
  STupleKey w1, w2;
  memset(&w1, 0, sizeof(w1));
  memset(&w2, 0, sizeof(w2));

  char* p1 = (char*)aBuf;
  char* p2 = (char*)bBuf;

  p1 = taosDecodeFixedU64(p1, &w1.groupId);
  p2 = taosDecodeFixedU64(p2, &w2.groupId);

  p1 = taosDecodeFixedI64(p1, &w1.ts);
  p2 = taosDecodeFixedI64(p2, &w2.ts);

  p1 = taosDecodeFixedI32(p1, &w1.exprIdx);
  p2 = taosDecodeFixedI32(p2, &w2.exprIdx);

  return STupleKeyCmpr(&w1, sizeof(w1), &w2, sizeof(w2));
}

int tupleKeyEncode(void* k, char* buf) {
  STupleKey* key = k;
  int        len = 0;
  len += taosEncodeFixedU64((void**)&buf, key->groupId);
  len += taosEncodeFixedI64((void**)&buf, key->ts);
  len += taosEncodeFixedI32((void**)&buf, key->exprIdx);
  return len;
}
int tupleKeyDecode(void* k, char* buf) {
  STupleKey* key = k;
  int        len = 0;
  char*      p = buf;
  p = taosDecodeFixedU64(p, &key->groupId);
  p = taosDecodeFixedI64(p, &key->ts);
  p = taosDecodeFixedI32(p, &key->exprIdx);
  return len;
}
dengyihao's avatar
dengyihao 已提交
589 590 591
int tupleKeyToString(void* k, char* buf) {
  int        n = 0;
  STupleKey* key = k;
dengyihao's avatar
dengyihao 已提交
592 593 594
  n += sprintf(buf + n, "[groupId:%" PRIu64 ",", key->groupId);
  n += sprintf(buf + n, "ts:%" PRIi64 ",", key->ts);
  n += sprintf(buf + n, "exprIdx:%d]", key->exprIdx);
dengyihao's avatar
dengyihao 已提交
595 596
  return n;
}
dengyihao's avatar
dengyihao 已提交
597 598 599 600 601 602 603 604 605 606 607 608 609 610 611 612 613 614 615 616 617 618 619 620 621 622 623 624

int parKeyDBComp(void* state, const char* aBuf, size_t aLen, const char* bBuf, size_t bLen) {
  int64_t w1, w2;
  memset(&w1, 0, sizeof(w1));
  memset(&w2, 0, sizeof(w2));
  char* p1 = (char*)aBuf;
  char* p2 = (char*)bBuf;

  taosDecodeFixedI64(p1, &w1);
  taosDecodeFixedI64(p2, &w2);
  if (w1 == w2) {
    return 0;
  } else {
    return w1 < w2 ? -1 : 1;
  }
}
int parKeyEncode(void* k, char* buf) {
  int64_t* groupid = k;
  int      len = taosEncodeFixedI64((void**)&buf, *groupid);
  return len;
}
int parKeyDecode(void* k, char* buf) {
  char*    p = buf;
  int64_t* groupid = k;

  p = taosDecodeFixedI64(p, groupid);
  return p - buf;
}
dengyihao's avatar
dengyihao 已提交
625 626 627
int parKeyToString(void* k, char* buf) {
  int64_t* key = k;
  int      n = 0;
dengyihao's avatar
dengyihao 已提交
628
  n = sprintf(buf + n, "[groupId:%" PRIi64 "]", *key);
dengyihao's avatar
dengyihao 已提交
629 630
  return n;
}
dengyihao's avatar
dengyihao 已提交
631
int stremaValueEncode(void* k, char* buf) {
dengyihao's avatar
dengyihao 已提交
632 633
  int           len = 0;
  SStreamValue* key = k;
dengyihao's avatar
dengyihao 已提交
634 635 636 637 638 639
  len += taosEncodeFixedI64((void**)&buf, key->unixTimestamp);
  len += taosEncodeFixedI32((void**)&buf, key->len);
  len += taosEncodeBinary((void**)&buf, key->data, key->len);
  return len;
}
int streamValueDecode(void* k, char* buf) {
dengyihao's avatar
dengyihao 已提交
640 641
  SStreamValue* key = k;
  char*         p = buf;
dengyihao's avatar
dengyihao 已提交
642 643 644 645 646 647
  p = taosDecodeFixedI64(p, &key->unixTimestamp);
  p = taosDecodeFixedI32(p, &key->len);
  p = taosDecodeBinary(p, (void**)&key->data, key->len);
  return p - buf;
}
int32_t streamValueToString(void* k, char* buf) {
dengyihao's avatar
dengyihao 已提交
648 649
  SStreamValue* key = k;
  int           n = 0;
dengyihao's avatar
dengyihao 已提交
650 651 652 653 654 655
  n += sprintf(buf + n, "[unixTimestamp:%" PRIi64 ",", key->unixTimestamp);
  n += sprintf(buf + n, "len:%d,", key->len);
  n += sprintf(buf + n, "data:%s]", key->data);
  return n;
}

dengyihao's avatar
dengyihao 已提交
656
/*1: stale, 0: no stale*/
dengyihao's avatar
dengyihao 已提交
657
int32_t streaValueIsStale(void* k, int64_t ts) {
dengyihao's avatar
dengyihao 已提交
658
  SStreamValue* key = k;
dengyihao's avatar
dengyihao 已提交
659 660 661 662 663
  if (key->unixTimestamp < ts) {
    return 1;
  }
  return 0;
}
dengyihao's avatar
dengyihao 已提交
664

dengyihao's avatar
dengyihao 已提交
665 666 667 668
void destroyFunc(void* arg) {
  (void)arg;
  return;
}
dengyihao's avatar
dengyihao 已提交
669

dengyihao's avatar
dengyihao 已提交
670 671
int32_t encodeValueFunc(void* value, int32_t vlen, int64_t ttl, char** dest) {
  SStreamValue key = {.unixTimestamp = ttl, .len = vlen, .data = (char*)(value)};
dengyihao's avatar
dengyihao 已提交
672 673 674 675 676 677 678 679 680 681 682 683 684 685
  int32_t      len = 0;
  if (*dest == NULL) {
    char* p = taosMemoryCalloc(1, sizeof(int64_t) + sizeof(int32_t) + key.len);
    char* buf = p;
    len += taosEncodeFixedI64((void**)&buf, key.unixTimestamp);
    len += taosEncodeFixedI32((void**)&buf, key.len);
    len += taosEncodeBinary((void**)&buf, (char*)value, vlen);
    *dest = p;
  } else {
    char* buf = *dest;
    len += taosEncodeFixedI64((void**)&buf, key.unixTimestamp);
    len += taosEncodeFixedI32((void**)&buf, key.len);
    len += taosEncodeBinary((void**)&buf, (char*)value, vlen);
  }
dengyihao's avatar
dengyihao 已提交
686 687
  return len;
}
dengyihao's avatar
dengyihao 已提交
688 689 690 691
/*
 *  ret >= 0 : found valid value
 *  ret < 0 : error or timeout
 */
dengyihao's avatar
dengyihao 已提交
692 693 694
int32_t decodeValueFunc(void* value, int32_t vlen, int64_t* ttl, char** dest) {
  SStreamValue key = {0};
  char*        p = value;
dengyihao's avatar
dengyihao 已提交
695
  if (streamStateValueIsStale(p)) {
dengyihao's avatar
dengyihao 已提交
696 697 698
    *dest = NULL;
    return -1;
  }
dengyihao's avatar
dengyihao 已提交
699 700
  p = taosDecodeFixedI64(p, &key.unixTimestamp);
  p = taosDecodeFixedI32(p, &key.len);
Y
yihaoDeng 已提交
701
  if (vlen != (sizeof(int64_t) + sizeof(int32_t) + key.len)) {
dengyihao's avatar
dengyihao 已提交
702
    if (dest != NULL) *dest = NULL;
Y
yihaoDeng 已提交
703
    qError("vlen: %d, read len: %d", vlen, key.len);
dengyihao's avatar
dengyihao 已提交
704 705 706
    return -1;
  }

dengyihao's avatar
dengyihao 已提交
707 708 709 710 711 712
  if (key.len == 0) {
    key.data = NULL;
  } else {
    p = taosDecodeBinary(p, (void**)&(key.data), key.len);
  }

dengyihao's avatar
dengyihao 已提交
713
  if (ttl != NULL) {
dengyihao's avatar
dengyihao 已提交
714
    int64_t now = taosGetTimestampMs();
dengyihao's avatar
dengyihao 已提交
715 716 717 718 719 720 721
    *ttl = key.unixTimestamp == 0 ? 0 : key.unixTimestamp - now;
  }
  if (dest != NULL) {
    *dest = key.data;
  } else {
    taosMemoryFree(key.data);
  }
dengyihao's avatar
dengyihao 已提交
722 723
  return key.len;
}
dengyihao's avatar
dengyihao 已提交
724

dengyihao's avatar
dengyihao 已提交
725 726 727 728 729 730 731 732 733 734 735 736 737 738 739 740 741 742 743 744 745 746 747 748 749 750 751 752
const char* compareDefaultName(void* arg) {
  (void)arg;
  return ginitDict[0].key;
}
const char* compareStateName(void* arg) {
  (void)arg;
  return ginitDict[1].key;
}
const char* compareWinKeyName(void* arg) {
  (void)arg;
  return ginitDict[2].key;
}
const char* compareSessionKeyName(void* arg) {
  (void)arg;
  return ginitDict[3].key;
}
const char* compareFuncKeyName(void* arg) {
  (void)arg;
  return ginitDict[4].key;
}
const char* compareParKeyName(void* arg) {
  (void)arg;
  return ginitDict[5].key;
}
const char* comparePartagKeyName(void* arg) {
  (void)arg;
  return ginitDict[6].key;
}
dengyihao's avatar
dengyihao 已提交
753

dengyihao's avatar
dengyihao 已提交
754 755 756 757 758 759 760 761
void destroyCompactFilteFactory(void* arg) {
  if (arg == NULL) return;
}
const char* compactFilteFactoryName(void* arg) {
  SCompactFilteFactory* state = arg;
  return "stream_compact_filter";
}

dengyihao's avatar
dengyihao 已提交
762
void          destroyCompactFilte(void* arg) { (void)arg; }
dengyihao's avatar
dengyihao 已提交
763 764
unsigned char compactFilte(void* arg, int level, const char* key, size_t klen, const char* val, size_t vlen,
                           char** newval, size_t* newvlen, unsigned char* value_changed) {
dengyihao's avatar
dengyihao 已提交
765
  return streamStateValueIsStale((char*)val) ? 1 : 0;
dengyihao's avatar
dengyihao 已提交
766 767 768 769 770 771 772 773 774
}
const char* compactFilteName(void* arg) { return "stream_filte"; }

rocksdb_compactionfilter_t* compactFilteFactoryCreateFilter(void* arg, rocksdb_compactionfiltercontext_t* ctx) {
  SCompactFilteFactory*       state = arg;
  rocksdb_compactionfilter_t* filter =
      rocksdb_compactionfilter_create(NULL, destroyCompactFilte, compactFilte, compactFilteName);
  return filter;
}
dengyihao's avatar
dengyihao 已提交
775

dengyihao's avatar
dengyihao 已提交
776 777
void destroyRocksdbCfInst(RocksdbCfInst* inst) {
  int cfLen = sizeof(ginitDict) / sizeof(ginitDict[0]);
dengyihao's avatar
dengyihao 已提交
778
  for (int i = 0; i < cfLen; i++) {
Y
yihaoDeng 已提交
779
    if (inst->pHandle[i]) rocksdb_column_family_handle_destroy((inst->pHandle)[i]);
dengyihao's avatar
dengyihao 已提交
780 781
  }

dengyihao's avatar
dengyihao 已提交
782 783 784 785 786 787 788 789 790
  rocksdb_writeoptions_destroy(inst->wOpt);
  inst->wOpt = NULL;

  rocksdb_readoptions_destroy(inst->rOpt);
  taosMemoryFree(inst->cfOpt);
  taosMemoryFreeClear(inst->param);
  taosMemoryFree(inst);
}

Y
yihaoDeng 已提交
791
int32_t streamStateOpenBackendCf(void* backend, char* name, char** cfs, int32_t nCf) {
Y
yihaoDeng 已提交
792 793 794 795 796
  SBackendWrapper* handle = backend;
  char*            err = NULL;
  int64_t          streamId;
  int32_t          taskId, dummy = 0;
  char             suffix[64] = {0};
Y
yihaoDeng 已提交
797 798 799 800 801 802 803 804 805

  rocksdb_options_t**              cfOpts = taosMemoryCalloc(nCf, sizeof(rocksdb_options_t*));
  RocksdbCfParam*                  params = taosMemoryCalloc(nCf, sizeof(RocksdbCfParam*));
  rocksdb_comparator_t**           pCompare = taosMemoryCalloc(nCf, sizeof(rocksdb_comparator_t**));
  rocksdb_column_family_handle_t** cfHandle = taosMemoryCalloc(nCf, sizeof(rocksdb_column_family_handle_t*));

  for (int i = 0; i < nCf; i++) {
    char* cf = cfs[i];
    char  funcname[64] = {0};
dengyihao's avatar
dengyihao 已提交
806
    cfOpts[i] = rocksdb_options_create_copy(handle->dbOpt);
Y
yihaoDeng 已提交
807 808 809 810
    if (i == 0) continue;
    if (3 == sscanf(cf, "0x%" PRIx64 "-%d_%s", &streamId, &taskId, funcname)) {
      rocksdb_block_based_table_options_t* tableOpt = rocksdb_block_based_options_create();
      rocksdb_block_based_options_set_block_cache(tableOpt, handle->cache);
dengyihao's avatar
dengyihao 已提交
811

Y
yihaoDeng 已提交
812 813
      rocksdb_filterpolicy_t* filter = rocksdb_filterpolicy_create_bloom(15);
      rocksdb_block_based_options_set_filter_policy(tableOpt, filter);
dengyihao's avatar
dengyihao 已提交
814

Y
yihaoDeng 已提交
815 816
      rocksdb_options_set_block_based_table_factory((rocksdb_options_t*)cfOpts[i], tableOpt);
      params[i].tableOpt = tableOpt;
dengyihao's avatar
dengyihao 已提交
817

dengyihao's avatar
dengyihao 已提交
818
      int      idx = streamStateGetCfIdx(NULL, funcname);
Y
yihaoDeng 已提交
819
      SCfInit* cfPara = &ginitDict[idx];
dengyihao's avatar
dengyihao 已提交
820

Y
yihaoDeng 已提交
821 822 823 824 825
      rocksdb_comparator_t* compare =
          rocksdb_comparator_create(NULL, cfPara->detroyFunc, cfPara->cmpFunc, cfPara->cmpName);
      rocksdb_options_set_comparator((rocksdb_options_t*)cfOpts[i], compare);
      pCompare[i] = compare;
    }
dengyihao's avatar
dengyihao 已提交
826
  }
Y
yihaoDeng 已提交
827
  rocksdb_t* db = rocksdb_open_column_families(handle->dbOpt, name, nCf, (const char* const*)cfs,
dengyihao's avatar
dengyihao 已提交
828
                                               (const rocksdb_options_t* const*)cfOpts, cfHandle, &err);
dengyihao's avatar
dengyihao 已提交
829 830 831 832
  if (err != NULL) {
    qError("failed to open rocksdb cf, reason:%s", err);
    taosMemoryFree(err);
  } else {
Y
yihaoDeng 已提交
833
    qDebug("succ to open rocksdb cf");
dengyihao's avatar
dengyihao 已提交
834
  }
Y
yihaoDeng 已提交
835
  // close default cf
836
  if (((rocksdb_column_family_handle_t**)cfHandle)[0] != 0) rocksdb_column_family_handle_destroy(cfHandle[0]);
Y
yihaoDeng 已提交
837 838
  rocksdb_options_destroy(cfOpts[0]);
  handle->db = db;
dengyihao's avatar
dengyihao 已提交
839

Y
yihaoDeng 已提交
840 841 842 843 844 845 846 847 848
  static int32_t cfLen = sizeof(ginitDict) / sizeof(ginitDict[0]);
  for (int i = 0; i < nCf; i++) {
    char* cf = cfs[i];
    if (i == 0) continue;
    char funcname[64] = {0};
    if (3 == sscanf(cf, "0x%" PRIx64 "-%d_%s", &streamId, &taskId, funcname)) {
      char idstr[128] = {0};
      sprintf(idstr, "0x%" PRIx64 "-%d", streamId, taskId);

dengyihao's avatar
dengyihao 已提交
849
      int idx = streamStateGetCfIdx(NULL, funcname);
Y
yihaoDeng 已提交
850 851

      RocksdbCfInst*  inst = NULL;
Y
yihaoDeng 已提交
852
      RocksdbCfInst** pInst = taosHashGet(handle->cfInst, idstr, strlen(idstr) + 1);
Y
yihaoDeng 已提交
853 854 855 856 857 858 859 860 861 862 863 864 865 866 867 868 869 870 871 872 873
      if (pInst == NULL || *pInst == NULL) {
        inst = taosMemoryCalloc(1, sizeof(RocksdbCfInst));
        inst->pHandle = taosMemoryCalloc(cfLen, sizeof(rocksdb_column_family_handle_t*));
        inst->cfOpt = taosMemoryCalloc(cfLen, sizeof(rocksdb_options_t*));
        inst->wOpt = rocksdb_writeoptions_create();
        inst->rOpt = rocksdb_readoptions_create();
        inst->param = taosMemoryCalloc(cfLen, sizeof(RocksdbCfParam));
        inst->pBackend = handle;
        inst->db = db;
        inst->pCompares = taosMemoryCalloc(cfLen, sizeof(rocksdb_comparator_t*));

        inst->dbOpt = handle->dbOpt;
        rocksdb_writeoptions_disable_WAL(inst->wOpt, 1);
        taosHashPut(handle->cfInst, idstr, strlen(idstr) + 1, &inst, sizeof(void*));
      } else {
        inst = *pInst;
      }
      inst->cfOpt[idx] = cfOpts[i];
      inst->pCompares[idx] = pCompare[i];
      memcpy(&(inst->param[idx]), &(params[i]), sizeof(RocksdbCfParam));
      inst->pHandle[idx] = cfHandle[i];
dengyihao's avatar
dengyihao 已提交
874
    }
Y
yihaoDeng 已提交
875 876
  }
  void** pIter = taosHashIterate(handle->cfInst, NULL);
Y
yihaoDeng 已提交
877
  while (pIter) {
Y
yihaoDeng 已提交
878
    RocksdbCfInst* inst = *pIter;
879 880 881 882 883 884 885 886 887 888 889 890 891 892 893 894 895 896 897 898 899 900 901 902

    for (int i = 0; i < cfLen; i++) {
      if (inst->cfOpt[i] == NULL) {
        rocksdb_options_t*                   opt = rocksdb_options_create_copy(handle->dbOpt);
        rocksdb_block_based_table_options_t* tableOpt = rocksdb_block_based_options_create();
        rocksdb_block_based_options_set_block_cache(tableOpt, handle->cache);

        rocksdb_filterpolicy_t* filter = rocksdb_filterpolicy_create_bloom(15);
        rocksdb_block_based_options_set_filter_policy(tableOpt, filter);

        rocksdb_options_set_block_based_table_factory((rocksdb_options_t*)opt, tableOpt);

        SCfInit* cfPara = &ginitDict[i];

        rocksdb_comparator_t* compare =
            rocksdb_comparator_create(NULL, cfPara->detroyFunc, cfPara->cmpFunc, cfPara->cmpName);
        rocksdb_options_set_comparator((rocksdb_options_t*)opt, compare);

        inst->pCompares[i] = compare;
        inst->cfOpt[i] = opt;
        inst->param[i].tableOpt = tableOpt;
      }
    }
    SCfComparator compare = {.comp = inst->pCompares, .numOfComp = cfLen};
dengyihao's avatar
dengyihao 已提交
903
    inst->pCompareNode = streamBackendAddCompare(handle, &compare);
904
    pIter = taosHashIterate(handle->cfInst, pIter);
dengyihao's avatar
dengyihao 已提交
905
  }
dengyihao's avatar
dengyihao 已提交
906

dengyihao's avatar
dengyihao 已提交
907 908 909 910 911
  taosMemoryFree(cfHandle);
  taosMemoryFree(pCompare);
  taosMemoryFree(params);
  taosMemoryFree(cfOpts);
  return 0;
dengyihao's avatar
dengyihao 已提交
912
}
dengyihao's avatar
dengyihao 已提交
913
int streamStateOpenBackend(void* backend, SStreamState* pState) {
dengyihao's avatar
dengyihao 已提交
914 915
  qInfo("start to open state %p on backend %p 0x%" PRIx64 "-%d", pState, backend, pState->streamId, pState->taskId);
  taosAcquireRef(streamBackendId, pState->streamBackendRid);
Y
yihaoDeng 已提交
916 917
  SBackendWrapper*   handle = backend;
  SBackendCfWrapper* pBackendCfWrapper = taosMemoryCalloc(1, sizeof(SBackendCfWrapper));
dengyihao's avatar
dengyihao 已提交
918
  taosThreadMutexLock(&handle->cfMutex);
Y
yihaoDeng 已提交
919

dengyihao's avatar
dengyihao 已提交
920
  RocksdbCfInst** ppInst = taosHashGet(handle->cfInst, pState->pTdbState->idstr, strlen(pState->pTdbState->idstr) + 1);
dengyihao's avatar
dengyihao 已提交
921 922
  if (ppInst != NULL && *ppInst != NULL) {
    RocksdbCfInst* inst = *ppInst;
Y
yihaoDeng 已提交
923 924 925 926 927 928 929 930 931
    pBackendCfWrapper->rocksdb = inst->db;
    pBackendCfWrapper->pHandle = (void**)inst->pHandle;
    pBackendCfWrapper->writeOpts = inst->wOpt;
    pBackendCfWrapper->readOpts = inst->rOpt;
    pBackendCfWrapper->cfOpts = (void**)(inst->cfOpt);
    pBackendCfWrapper->dbOpt = handle->dbOpt;
    pBackendCfWrapper->param = inst->param;
    pBackendCfWrapper->pBackend = handle;
    pBackendCfWrapper->pComparNode = inst->pCompareNode;
dengyihao's avatar
dengyihao 已提交
932
    taosThreadMutexUnlock(&handle->cfMutex);
Y
yihaoDeng 已提交
933 934
    pBackendCfWrapper->backendId = pState->streamBackendRid;
    memcpy(pBackendCfWrapper->idstr, pState->pTdbState->idstr, sizeof(pState->pTdbState->idstr));
Y
yihaoDeng 已提交
935

Y
yihaoDeng 已提交
936 937 938 939
    int64_t id = taosAddRef(streamBackendCfWrapperId, pBackendCfWrapper);
    pState->pTdbState->backendCfWrapperId = id;
    pState->pTdbState->pBackendCfWrapper = pBackendCfWrapper;
    qInfo("succ to open state %p on backendWrapper, %p, %s", pState, pBackendCfWrapper, pBackendCfWrapper->idstr);
dengyihao's avatar
dengyihao 已提交
940 941 942 943
    return 0;
  }
  taosThreadMutexUnlock(&handle->cfMutex);

dengyihao's avatar
dengyihao 已提交
944
  char* err = NULL;
dengyihao's avatar
dengyihao 已提交
945
  int   cfLen = sizeof(ginitDict) / sizeof(ginitDict[0]);
dengyihao's avatar
dengyihao 已提交
946

dengyihao's avatar
dengyihao 已提交
947
  RocksdbCfParam*           param = taosMemoryCalloc(cfLen, sizeof(RocksdbCfParam));
dengyihao's avatar
dengyihao 已提交
948 949
  const rocksdb_options_t** cfOpt = taosMemoryCalloc(cfLen, sizeof(rocksdb_options_t*));
  for (int i = 0; i < cfLen; i++) {
950
    cfOpt[i] = rocksdb_options_create_copy(handle->dbOpt);
dengyihao's avatar
dengyihao 已提交
951
    // refactor later
dengyihao's avatar
dengyihao 已提交
952
    rocksdb_block_based_table_options_t* tableOpt = rocksdb_block_based_options_create();
dengyihao's avatar
dengyihao 已提交
953
    rocksdb_block_based_options_set_block_cache(tableOpt, handle->cache);
dengyihao's avatar
dengyihao 已提交
954

dengyihao's avatar
dengyihao 已提交
955
    rocksdb_filterpolicy_t* filter = rocksdb_filterpolicy_create_bloom(15);
dengyihao's avatar
dengyihao 已提交
956
    rocksdb_block_based_options_set_filter_policy(tableOpt, filter);
dengyihao's avatar
dengyihao 已提交
957 958

    rocksdb_options_set_block_based_table_factory((rocksdb_options_t*)cfOpt[i], tableOpt);
dengyihao's avatar
dengyihao 已提交
959

dengyihao's avatar
dengyihao 已提交
960
    param[i].tableOpt = tableOpt;
dengyihao's avatar
dengyihao 已提交
961
  };
dengyihao's avatar
dengyihao 已提交
962

dengyihao's avatar
dengyihao 已提交
963
  rocksdb_comparator_t** pCompare = taosMemoryCalloc(cfLen, sizeof(rocksdb_comparator_t**));
dengyihao's avatar
dengyihao 已提交
964
  for (int i = 0; i < cfLen; i++) {
dengyihao's avatar
dengyihao 已提交
965 966
    SCfInit* cf = &ginitDict[i];

dengyihao's avatar
dengyihao 已提交
967 968 969 970
    rocksdb_comparator_t* compare = rocksdb_comparator_create(NULL, cf->detroyFunc, cf->cmpFunc, cf->cmpName);
    rocksdb_options_set_comparator((rocksdb_options_t*)cfOpt[i], compare);
    pCompare[i] = compare;
  }
Y
yihaoDeng 已提交
971
  rocksdb_column_family_handle_t** cfHandle = taosMemoryCalloc(cfLen, sizeof(rocksdb_column_family_handle_t*));
Y
yihaoDeng 已提交
972 973 974 975 976 977 978 979 980 981
  pBackendCfWrapper->rocksdb = handle->db;
  pBackendCfWrapper->pHandle = (void**)cfHandle;
  pBackendCfWrapper->writeOpts = rocksdb_writeoptions_create();
  pBackendCfWrapper->readOpts = rocksdb_readoptions_create();
  pBackendCfWrapper->cfOpts = (void**)cfOpt;
  pBackendCfWrapper->dbOpt = handle->dbOpt;
  pBackendCfWrapper->param = param;
  pBackendCfWrapper->pBackend = handle;
  pBackendCfWrapper->backendId = pState->streamBackendRid;
  taosThreadRwlockInit(&pBackendCfWrapper->rwLock, NULL);
dengyihao's avatar
dengyihao 已提交
982
  SCfComparator compare = {.comp = pCompare, .numOfComp = cfLen};
Y
yihaoDeng 已提交
983 984 985 986 987 988 989 990
  pBackendCfWrapper->pComparNode = streamBackendAddCompare(handle, &compare);
  rocksdb_writeoptions_disable_WAL(pBackendCfWrapper->writeOpts, 1);
  memcpy(pBackendCfWrapper->idstr, pState->pTdbState->idstr, sizeof(pState->pTdbState->idstr));

  int64_t id = taosAddRef(streamBackendCfWrapperId, pBackendCfWrapper);
  pState->pTdbState->backendCfWrapperId = id;
  pState->pTdbState->pBackendCfWrapper = pBackendCfWrapper;
  qInfo("succ to open state %p on backendWrapper %p %s", pState, pBackendCfWrapper, pBackendCfWrapper->idstr);
dengyihao's avatar
dengyihao 已提交
991 992
  return 0;
}
993

dengyihao's avatar
dengyihao 已提交
994
void streamStateCloseBackend(SStreamState* pState, bool remove) {
Y
yihaoDeng 已提交
995 996
  SBackendCfWrapper* wrapper = pState->pTdbState->pBackendCfWrapper;
  SBackendWrapper*   pHandle = wrapper->pBackend;
dengyihao's avatar
dengyihao 已提交
997
  taosThreadMutexLock(&pHandle->cfMutex);
Y
yihaoDeng 已提交
998
  RocksdbCfInst** ppInst = taosHashGet(pHandle->cfInst, wrapper->idstr, strlen(pState->pTdbState->idstr) + 1);
dengyihao's avatar
dengyihao 已提交
999 1000 1001
  if (ppInst != NULL && *ppInst != NULL) {
    RocksdbCfInst* inst = *ppInst;
    taosMemoryFree(inst);
dengyihao's avatar
dengyihao 已提交
1002
    taosHashRemove(pHandle->cfInst, pState->pTdbState->idstr, strlen(pState->pTdbState->idstr) + 1);
dengyihao's avatar
dengyihao 已提交
1003 1004
  }
  taosThreadMutexUnlock(&pHandle->cfMutex);
dengyihao's avatar
dengyihao 已提交
1005

dengyihao's avatar
dengyihao 已提交
1006
  char* status[] = {"close", "drop"};
Y
yihaoDeng 已提交
1007 1008
  qInfo("start to close %s state %p on backendWrapper %p %s", status[remove == false ? 0 : 1], pState, wrapper,
        wrapper->idstr);
Y
yihaoDeng 已提交
1009
  wrapper->remove |= remove;  // update by other pState
Y
yihaoDeng 已提交
1010
  taosReleaseRef(streamBackendCfWrapperId, pState->pTdbState->backendCfWrapperId);
dengyihao's avatar
dengyihao 已提交
1011
}
dengyihao's avatar
dengyihao 已提交
1012 1013 1014
void streamStateDestroyCompar(void* arg) {
  SCfComparator* comp = (SCfComparator*)arg;
  for (int i = 0; i < comp->numOfComp; i++) {
1015
    if (comp->comp[i]) rocksdb_comparator_destroy(comp->comp[i]);
dengyihao's avatar
dengyihao 已提交
1016 1017 1018
  }
  taosMemoryFree(comp->comp);
}
1019

dengyihao's avatar
dengyihao 已提交
1020
int streamStateGetCfIdx(SStreamState* pState, const char* funcName) {
Y
yihaoDeng 已提交
1021
  int    idx = -1;
dengyihao's avatar
dengyihao 已提交
1022
  size_t len = strlen(funcName);
dengyihao's avatar
dengyihao 已提交
1023
  for (int i = 0; i < sizeof(ginitDict) / sizeof(ginitDict[0]); i++) {
dengyihao's avatar
dengyihao 已提交
1024
    if (len == ginitDict[i].len && strncmp(funcName, ginitDict[i].key, strlen(funcName)) == 0) {
Y
yihaoDeng 已提交
1025 1026
      idx = i;
      break;
dengyihao's avatar
dengyihao 已提交
1027 1028
    }
  }
1029
  if (pState != NULL && idx != -1) {
Y
yihaoDeng 已提交
1030
    SBackendCfWrapper*              wrapper = pState->pTdbState->pBackendCfWrapper;
Y
yihaoDeng 已提交
1031
    rocksdb_column_family_handle_t* cf = NULL;
Y
yihaoDeng 已提交
1032 1033 1034
    taosThreadRwlockRdlock(&wrapper->rwLock);
    cf = wrapper->pHandle[idx];
    taosThreadRwlockUnlock(&wrapper->rwLock);
Y
yihaoDeng 已提交
1035 1036
    if (cf == NULL) {
      char buf[128] = {0};
Y
yihaoDeng 已提交
1037
      GEN_COLUMN_FAMILY_NAME(buf, wrapper->idstr, ginitDict[idx].key);
Y
yihaoDeng 已提交
1038
      char* err = NULL;
1039

Y
yihaoDeng 已提交
1040 1041
      taosThreadRwlockWrlock(&wrapper->rwLock);
      cf = rocksdb_create_column_family(wrapper->rocksdb, wrapper->cfOpts[idx], buf, &err);
Y
yihaoDeng 已提交
1042 1043
      if (err != NULL) {
        idx = -1;
Y
yihaoDeng 已提交
1044
        qError("failed to to open cf, %p %s_%s, reason:%s", pState, wrapper->idstr, funcName, err);
Y
yihaoDeng 已提交
1045
        taosMemoryFree(err);
Y
yihaoDeng 已提交
1046
      } else {
Y
yihaoDeng 已提交
1047
        wrapper->pHandle[idx] = cf;
Y
yihaoDeng 已提交
1048
      }
Y
yihaoDeng 已提交
1049
      taosThreadRwlockUnlock(&wrapper->rwLock);
Y
yihaoDeng 已提交
1050 1051 1052 1053
    }
  }

  return idx;
dengyihao's avatar
dengyihao 已提交
1054
}
dengyihao's avatar
dengyihao 已提交
1055
bool streamStateIterSeekAndValid(rocksdb_iterator_t* iter, char* buf, size_t len) {
dengyihao's avatar
dengyihao 已提交
1056
  rocksdb_iter_seek(iter, buf, len);
dengyihao's avatar
dengyihao 已提交
1057
  if (!rocksdb_iter_valid(iter)) {
dengyihao's avatar
dengyihao 已提交
1058 1059 1060 1061
    rocksdb_iter_seek_for_prev(iter, buf, len);
    if (!rocksdb_iter_valid(iter)) {
      return false;
    }
dengyihao's avatar
dengyihao 已提交
1062
  }
dengyihao's avatar
dengyihao 已提交
1063
  return true;
dengyihao's avatar
dengyihao 已提交
1064
}
dengyihao's avatar
dengyihao 已提交
1065 1066
rocksdb_iterator_t* streamStateIterCreate(SStreamState* pState, const char* cfName, rocksdb_snapshot_t** snapshot,
                                          rocksdb_readoptions_t** readOpt) {
dengyihao's avatar
dengyihao 已提交
1067
  int idx = streamStateGetCfIdx(pState, cfName);
dengyihao's avatar
dengyihao 已提交
1068

Y
yihaoDeng 已提交
1069
  SBackendCfWrapper* wrapper = pState->pTdbState->pBackendCfWrapper;
dengyihao's avatar
dengyihao 已提交
1070
  if (snapshot != NULL) {
Y
yihaoDeng 已提交
1071
    *snapshot = (rocksdb_snapshot_t*)rocksdb_create_snapshot(wrapper->rocksdb);
dengyihao's avatar
dengyihao 已提交
1072
  }
dengyihao's avatar
dengyihao 已提交
1073 1074
  rocksdb_readoptions_t* rOpt = rocksdb_readoptions_create();
  *readOpt = rOpt;
dengyihao's avatar
dengyihao 已提交
1075

dengyihao's avatar
dengyihao 已提交
1076
  rocksdb_readoptions_set_snapshot(rOpt, *snapshot);
dengyihao's avatar
dengyihao 已提交
1077 1078
  rocksdb_readoptions_set_fill_cache(rOpt, 0);

Y
yihaoDeng 已提交
1079
  return rocksdb_create_iterator_cf(wrapper->rocksdb, rOpt, ((rocksdb_column_family_handle_t**)wrapper->pHandle)[idx]);
dengyihao's avatar
dengyihao 已提交
1080 1081
}

Y
yihaoDeng 已提交
1082 1083 1084 1085 1086
#define STREAM_STATE_PUT_ROCKSDB(pState, funcname, key, value, vLen)                                                   \
  do {                                                                                                                 \
    code = 0;                                                                                                          \
    char  buf[128] = {0};                                                                                              \
    char* err = NULL;                                                                                                  \
dengyihao's avatar
dengyihao 已提交
1087
    int   i = streamStateGetCfIdx(pState, funcname);                                                                   \
Y
yihaoDeng 已提交
1088 1089 1090 1091 1092
    if (i < 0) {                                                                                                       \
      qWarn("streamState failed to get cf name: %s", funcname);                                                        \
      code = -1;                                                                                                       \
      break;                                                                                                           \
    }                                                                                                                  \
Y
yihaoDeng 已提交
1093 1094
    SBackendCfWrapper* wrapper = pState->pTdbState->pBackendCfWrapper;                                                 \
    char               toString[128] = {0};                                                                            \
Y
yihaoDeng 已提交
1095 1096
    if (qDebugFlag & DEBUG_TRACE) ginitDict[i].toStrFunc((void*)key, toString);                                        \
    int32_t                         klen = ginitDict[i].enFunc((void*)key, buf);                                       \
Y
yihaoDeng 已提交
1097 1098 1099 1100 1101
    rocksdb_column_family_handle_t* pHandle = ((rocksdb_column_family_handle_t**)wrapper->pHandle)[ginitDict[i].idx];  \
    rocksdb_t*                      db = wrapper->rocksdb;                                                             \
    rocksdb_writeoptions_t*         opts = wrapper->writeOpts;                                                         \
    char*                           ttlV = NULL;                                                                       \
    int32_t                         ttlVLen = ginitDict[i].enValueFunc((char*)value, vLen, 0, &ttlV);                  \
Y
yihaoDeng 已提交
1102 1103 1104 1105 1106 1107 1108 1109 1110
    rocksdb_put_cf(db, opts, pHandle, (const char*)buf, klen, (const char*)ttlV, (size_t)ttlVLen, &err);               \
    if (err != NULL) {                                                                                                 \
      taosMemoryFree(err);                                                                                             \
      qError("streamState str: %s failed to write to %s, err: %s", toString, funcname, err);                           \
      code = -1;                                                                                                       \
    } else {                                                                                                           \
      qTrace("streamState str:%s succ to write to %s, rowValLen:%d, ttlValLen:%d", toString, funcname, vLen, ttlVLen); \
    }                                                                                                                  \
    taosMemoryFree(ttlV);                                                                                              \
dengyihao's avatar
dengyihao 已提交
1111 1112
  } while (0);

Y
yihaoDeng 已提交
1113 1114 1115 1116 1117 1118 1119 1120 1121 1122 1123
#define STREAM_STATE_GET_ROCKSDB(pState, funcname, key, pVal, vLen)                                                   \
  do {                                                                                                                \
    code = 0;                                                                                                         \
    char  buf[128] = {0};                                                                                             \
    char* err = NULL;                                                                                                 \
    int   i = streamStateGetCfIdx(pState, funcname);                                                                  \
    if (i < 0) {                                                                                                      \
      qWarn("streamState failed to get cf name: %s", funcname);                                                       \
      code = -1;                                                                                                      \
      break;                                                                                                          \
    }                                                                                                                 \
Y
yihaoDeng 已提交
1124 1125
    SBackendCfWrapper* wrapper = pState->pTdbState->pBackendCfWrapper;                                                \
    char               toString[128] = {0};                                                                           \
Y
yihaoDeng 已提交
1126 1127 1128 1129 1130 1131 1132 1133 1134 1135 1136 1137 1138 1139 1140 1141 1142 1143 1144 1145 1146 1147 1148 1149 1150 1151 1152 1153 1154
    if (qDebugFlag & DEBUG_TRACE) ginitDict[i].toStrFunc((void*)key, toString);                                       \
    int32_t                         klen = ginitDict[i].enFunc((void*)key, buf);                                      \
    rocksdb_column_family_handle_t* pHandle = ((rocksdb_column_family_handle_t**)wrapper->pHandle)[ginitDict[i].idx]; \
    rocksdb_t*                      db = wrapper->rocksdb;                                                            \
    rocksdb_readoptions_t*          opts = wrapper->readOpts;                                                         \
    size_t                          len = 0;                                                                          \
    char* val = rocksdb_get_cf(db, opts, pHandle, (const char*)buf, klen, (size_t*)&len, &err);                       \
    if (val == NULL || len == 0) {                                                                                    \
      if (err == NULL) {                                                                                              \
        qTrace("streamState str: %s failed to read from %s_%s, err: not exist", toString, wrapper->idstr, funcname);  \
      } else {                                                                                                        \
        qError("streamState str: %s failed to read from %s_%s, err: %s", toString, wrapper->idstr, funcname, err);    \
        taosMemoryFreeClear(err);                                                                                     \
      }                                                                                                               \
      code = -1;                                                                                                      \
    } else {                                                                                                          \
      char*   p = NULL;                                                                                               \
      int32_t tlen = ginitDict[i].deValueFunc(val, len, NULL, (char**)pVal);                                          \
      if (tlen <= 0) {                                                                                                \
        qError("streamState str: %s failed to read from %s_%s, err: already ttl ", toString, wrapper->idstr,          \
               funcname);                                                                                             \
        code = -1;                                                                                                    \
      } else {                                                                                                        \
        qTrace("streamState str: %s succ to read from %s_%s, valLen:%d", toString, wrapper->idstr, funcname, tlen);   \
      }                                                                                                               \
      taosMemoryFree(val);                                                                                            \
      if (vLen != NULL) *vLen = tlen;                                                                                 \
    }                                                                                                                 \
    if (code == 0) qDebug("streamState str: %s succ to read from %s_%s", toString, wrapper->idstr, funcname);         \
dengyihao's avatar
dengyihao 已提交
1155 1156
  } while (0);

Y
yihaoDeng 已提交
1157 1158 1159 1160 1161 1162 1163 1164 1165 1166 1167
#define STREAM_STATE_DEL_ROCKSDB(pState, funcname, key)                                                               \
  do {                                                                                                                \
    code = 0;                                                                                                         \
    char  buf[128] = {0};                                                                                             \
    char* err = NULL;                                                                                                 \
    int   i = streamStateGetCfIdx(pState, funcname);                                                                  \
    if (i < 0) {                                                                                                      \
      qWarn("streamState failed to get cf name: %s_%s", pState->pTdbState->idstr, funcname);                          \
      code = -1;                                                                                                      \
      break;                                                                                                          \
    }                                                                                                                 \
Y
yihaoDeng 已提交
1168 1169
    SBackendCfWrapper* wrapper = pState->pTdbState->pBackendCfWrapper;                                                \
    char               toString[128] = {0};                                                                           \
Y
yihaoDeng 已提交
1170 1171 1172 1173 1174 1175 1176 1177 1178 1179 1180 1181 1182
    if (qDebugFlag & DEBUG_TRACE) ginitDict[i].toStrFunc((void*)key, toString);                                       \
    int32_t                         klen = ginitDict[i].enFunc((void*)key, buf);                                      \
    rocksdb_column_family_handle_t* pHandle = ((rocksdb_column_family_handle_t**)wrapper->pHandle)[ginitDict[i].idx]; \
    rocksdb_t*                      db = wrapper->rocksdb;                                                            \
    rocksdb_writeoptions_t*         opts = wrapper->writeOpts;                                                        \
    rocksdb_delete_cf(db, opts, pHandle, (const char*)buf, klen, &err);                                               \
    if (err != NULL) {                                                                                                \
      qError("streamState str: %s failed to del from %s_%s, err: %s", toString, wrapper->idstr, funcname, err);       \
      taosMemoryFree(err);                                                                                            \
      code = -1;                                                                                                      \
    } else {                                                                                                          \
      qTrace("streamState str: %s succ to del from %s_%s", toString, wrapper->idstr, funcname);                       \
    }                                                                                                                 \
dengyihao's avatar
dengyihao 已提交
1183 1184
  } while (0);

dengyihao's avatar
dengyihao 已提交
1185
// state cf
dengyihao's avatar
dengyihao 已提交
1186 1187 1188 1189
int32_t streamStatePut_rocksdb(SStreamState* pState, const SWinKey* key, const void* value, int32_t vLen) {
  int code = 0;

  SStateKey sKey = {.key = *key, .opNum = pState->number};
dengyihao's avatar
dengyihao 已提交
1190
  STREAM_STATE_PUT_ROCKSDB(pState, "state", &sKey, (void*)value, vLen);
dengyihao's avatar
dengyihao 已提交
1191 1192 1193 1194 1195
  return code;
}
int32_t streamStateGet_rocksdb(SStreamState* pState, const SWinKey* key, void** pVal, int32_t* pVLen) {
  int       code = 0;
  SStateKey sKey = {.key = *key, .opNum = pState->number};
dengyihao's avatar
dengyihao 已提交
1196
  STREAM_STATE_GET_ROCKSDB(pState, "state", &sKey, pVal, pVLen);
dengyihao's avatar
dengyihao 已提交
1197 1198 1199 1200 1201
  return code;
}
int32_t streamStateDel_rocksdb(SStreamState* pState, const SWinKey* key) {
  int       code = 0;
  SStateKey sKey = {.key = *key, .opNum = pState->number};
dengyihao's avatar
dengyihao 已提交
1202
  STREAM_STATE_DEL_ROCKSDB(pState, "state", &sKey);
dengyihao's avatar
dengyihao 已提交
1203 1204
  return code;
}
dengyihao's avatar
dengyihao 已提交
1205
int32_t streamStateClear_rocksdb(SStreamState* pState) {
dengyihao's avatar
dengyihao 已提交
1206
  qDebug("streamStateClear_rocksdb");
dengyihao's avatar
dengyihao 已提交
1207

Y
yihaoDeng 已提交
1208 1209 1210 1211 1212
  SBackendCfWrapper* wrapper = pState->pTdbState->pBackendCfWrapper;
  char               sKeyStr[128] = {0};
  char               eKeyStr[128] = {0};
  SStateKey          sKey = {.key = {.ts = 0, .groupId = 0}, .opNum = pState->number};
  SStateKey          eKey = {.key = {.ts = INT64_MAX, .groupId = UINT64_MAX}, .opNum = pState->number};
dengyihao's avatar
dengyihao 已提交
1213

dengyihao's avatar
dengyihao 已提交
1214 1215
  int sLen = stateKeyEncode(&sKey, sKeyStr);
  int eLen = stateKeyEncode(&eKey, eKeyStr);
dengyihao's avatar
dengyihao 已提交
1216

Y
yihaoDeng 已提交
1217
  if (wrapper->pHandle[1] != NULL) {
dengyihao's avatar
dengyihao 已提交
1218
    char* err = NULL;
Y
yihaoDeng 已提交
1219 1220
    rocksdb_delete_range_cf(wrapper->rocksdb, wrapper->writeOpts, wrapper->pHandle[1], sKeyStr, sLen, eKeyStr, eLen,
                            &err);
dengyihao's avatar
dengyihao 已提交
1221 1222 1223 1224 1225 1226 1227 1228 1229
    if (err != NULL) {
      char toStringStart[128] = {0};
      char toStringEnd[128] = {0};
      stateKeyToString(&sKey, toStringStart);
      stateKeyToString(&eKey, toStringEnd);

      qWarn("failed to delete range cf(state) start: %s, end:%s, reason:%s", toStringStart, toStringEnd, err);
      taosMemoryFree(err);
    } else {
Y
yihaoDeng 已提交
1230
      rocksdb_compact_range_cf(wrapper->rocksdb, wrapper->pHandle[1], sKeyStr, sLen, eKeyStr, eLen);
dengyihao's avatar
dengyihao 已提交
1231
    }
dengyihao's avatar
dengyihao 已提交
1232 1233
  }

dengyihao's avatar
dengyihao 已提交
1234
  return 0;
dengyihao's avatar
dengyihao 已提交
1235
}
dengyihao's avatar
dengyihao 已提交
1236 1237
int32_t streamStateCurNext_rocksdb(SStreamState* pState, SStreamStateCur* pCur) {
  if (!pCur) {
dengyihao's avatar
dengyihao 已提交
1238 1239
    return -1;
  }
dengyihao's avatar
dengyihao 已提交
1240
  rocksdb_iter_next(pCur->iter);
dengyihao's avatar
dengyihao 已提交
1241 1242
  return 0;
}
dengyihao's avatar
dengyihao 已提交
1243 1244 1245 1246 1247 1248 1249 1250
int32_t streamStateGetFirst_rocksdb(SStreamState* pState, SWinKey* key) {
  qDebug("streamStateGetFirst_rocksdb");
  SWinKey tmp = {.ts = 0, .groupId = 0};
  streamStatePut_rocksdb(pState, &tmp, NULL, 0);
  SStreamStateCur* pCur = streamStateSeekKeyNext_rocksdb(pState, &tmp);
  int32_t          code = streamStateGetKVByCur_rocksdb(pCur, key, NULL, 0);
  streamStateFreeCur(pCur);
  streamStateDel_rocksdb(pState, &tmp);
dengyihao's avatar
dengyihao 已提交
1251 1252 1253
  return code;
}

dengyihao's avatar
dengyihao 已提交
1254 1255 1256 1257
int32_t streamStateGetGroupKVByCur_rocksdb(SStreamStateCur* pCur, SWinKey* pKey, const void** pVal, int32_t* pVLen) {
  qDebug("streamStateGetGroupKVByCur_rocksdb");
  if (!pCur) {
    return -1;
dengyihao's avatar
dengyihao 已提交
1258
  }
dengyihao's avatar
dengyihao 已提交
1259
  uint64_t groupId = pKey->groupId;
dengyihao's avatar
dengyihao 已提交
1260

dengyihao's avatar
dengyihao 已提交
1261 1262 1263 1264 1265
  int32_t code = streamStateFillGetKVByCur_rocksdb(pCur, pKey, pVal, pVLen);
  if (code == 0) {
    if (pKey->groupId == groupId) {
      return 0;
    }
dengyihao's avatar
dengyihao 已提交
1266
  }
dengyihao's avatar
dengyihao 已提交
1267
  return -1;
dengyihao's avatar
dengyihao 已提交
1268
}
dengyihao's avatar
dengyihao 已提交
1269 1270 1271 1272 1273
int32_t streamStateAddIfNotExist_rocksdb(SStreamState* pState, const SWinKey* key, void** pVal, int32_t* pVLen) {
  qDebug("streamStateAddIfNotExist_rocksdb");
  int32_t size = *pVLen;
  if (streamStateGet_rocksdb(pState, key, pVal, pVLen) == 0) {
    return 0;
dengyihao's avatar
dengyihao 已提交
1274
  }
dengyihao's avatar
dengyihao 已提交
1275 1276 1277
  *pVal = taosMemoryMalloc(size);
  memset(*pVal, 0, size);
  return 0;
dengyihao's avatar
dengyihao 已提交
1278
}
dengyihao's avatar
dengyihao 已提交
1279 1280 1281
int32_t streamStateCurPrev_rocksdb(SStreamState* pState, SStreamStateCur* pCur) {
  qDebug("streamStateCurPrev_rocksdb");
  if (!pCur) return -1;
dengyihao's avatar
dengyihao 已提交
1282

dengyihao's avatar
dengyihao 已提交
1283 1284
  rocksdb_iter_prev(pCur->iter);
  return 0;
dengyihao's avatar
dengyihao 已提交
1285
}
dengyihao's avatar
dengyihao 已提交
1286 1287 1288 1289 1290
int32_t streamStateGetKVByCur_rocksdb(SStreamStateCur* pCur, SWinKey* pKey, const void** pVal, int32_t* pVLen) {
  qDebug("streamStateGetKVByCur_rocksdb");
  if (!pCur) return -1;
  SStateKey  tkey;
  SStateKey* pKtmp = &tkey;
dengyihao's avatar
dengyihao 已提交
1291

dengyihao's avatar
dengyihao 已提交
1292
  if (rocksdb_iter_valid(pCur->iter) && !iterValueIsStale(pCur->iter)) {
dengyihao's avatar
dengyihao 已提交
1293 1294 1295 1296 1297 1298 1299 1300 1301 1302 1303
    size_t tlen;
    char*  keyStr = (char*)rocksdb_iter_key(pCur->iter, &tlen);
    stateKeyDecode((void*)pKtmp, keyStr);
    if (pKtmp->opNum != pCur->number) {
      return -1;
    }
    size_t vlen = 0;
    if (pVal != NULL) *pVal = (char*)rocksdb_iter_value(pCur->iter, &vlen);
    if (pVLen != NULL) *pVLen = vlen;
    *pKey = pKtmp->key;
    return 0;
dengyihao's avatar
dengyihao 已提交
1304
  }
dengyihao's avatar
dengyihao 已提交
1305
  return -1;
dengyihao's avatar
dengyihao 已提交
1306
}
dengyihao's avatar
dengyihao 已提交
1307 1308 1309 1310 1311 1312
SStreamStateCur* streamStateGetAndCheckCur_rocksdb(SStreamState* pState, SWinKey* key) {
  qDebug("streamStateGetAndCheckCur_rocksdb");
  SStreamStateCur* pCur = streamStateFillGetCur_rocksdb(pState, key);
  if (pCur) {
    int32_t code = streamStateGetGroupKVByCur_rocksdb(pCur, key, NULL, 0);
    if (code == 0) return pCur;
dengyihao's avatar
dengyihao 已提交
1313 1314
    streamStateFreeCur(pCur);
  }
dengyihao's avatar
dengyihao 已提交
1315
  return NULL;
dengyihao's avatar
dengyihao 已提交
1316
}
dengyihao's avatar
dengyihao 已提交
1317

dengyihao's avatar
dengyihao 已提交
1318 1319
SStreamStateCur* streamStateSeekKeyNext_rocksdb(SStreamState* pState, const SWinKey* key) {
  qDebug("streamStateSeekKeyNext_rocksdb");
dengyihao's avatar
dengyihao 已提交
1320 1321 1322 1323
  SStreamStateCur* pCur = taosMemoryCalloc(1, sizeof(SStreamStateCur));
  if (pCur == NULL) {
    return NULL;
  }
Y
yihaoDeng 已提交
1324
  SBackendCfWrapper* wrapper = pState->pTdbState->pBackendCfWrapper;
dengyihao's avatar
dengyihao 已提交
1325
  pCur->number = pState->number;
Y
yihaoDeng 已提交
1326
  pCur->db = wrapper->rocksdb;
1327 1328
  pCur->iter = streamStateIterCreate(pState, "state", (rocksdb_snapshot_t**)&pCur->snapshot,
                                     (rocksdb_readoptions_t**)&pCur->readOpt);
dengyihao's avatar
dengyihao 已提交
1329

dengyihao's avatar
dengyihao 已提交
1330 1331 1332
  SStateKey sKey = {.key = *key, .opNum = pState->number};
  char      buf[128] = {0};
  int       len = stateKeyEncode((void*)&sKey, buf);
dengyihao's avatar
dengyihao 已提交
1333
  if (!streamStateIterSeekAndValid(pCur->iter, buf, len)) {
dengyihao's avatar
dengyihao 已提交
1334 1335
    streamStateFreeCur(pCur);
    return NULL;
dengyihao's avatar
dengyihao 已提交
1336
  }
dengyihao's avatar
dengyihao 已提交
1337
  // skip ttl expired data
dengyihao's avatar
dengyihao 已提交
1338
  while (rocksdb_iter_valid(pCur->iter) && iterValueIsStale(pCur->iter)) {
dengyihao's avatar
dengyihao 已提交
1339
    rocksdb_iter_next(pCur->iter);
dengyihao's avatar
dengyihao 已提交
1340
  }
dengyihao's avatar
dengyihao 已提交
1341

dengyihao's avatar
dengyihao 已提交
1342 1343 1344 1345 1346 1347 1348
  if (rocksdb_iter_valid(pCur->iter)) {
    SStateKey curKey;
    size_t    kLen;
    char*     keyStr = (char*)rocksdb_iter_key(pCur->iter, &kLen);
    stateKeyDecode((void*)&curKey, keyStr);
    if (stateKeyCmpr(&sKey, sizeof(sKey), &curKey, sizeof(curKey)) > 0) {
      return pCur;
dengyihao's avatar
dengyihao 已提交
1349
    }
dengyihao's avatar
dengyihao 已提交
1350 1351
    rocksdb_iter_next(pCur->iter);
    return pCur;
dengyihao's avatar
dengyihao 已提交
1352
  }
dengyihao's avatar
dengyihao 已提交
1353 1354
  streamStateFreeCur(pCur);
  return NULL;
dengyihao's avatar
dengyihao 已提交
1355 1356
}

dengyihao's avatar
dengyihao 已提交
1357 1358
SStreamStateCur* streamStateSeekToLast_rocksdb(SStreamState* pState, const SWinKey* key) {
  qDebug("streamStateGetCur_rocksdb");
Y
yihaoDeng 已提交
1359 1360
  int32_t            code = 0;
  SBackendCfWrapper* wrapper = pState->pTdbState->pBackendCfWrapper;
Y
yihaoDeng 已提交
1361

dengyihao's avatar
dengyihao 已提交
1362 1363
  const SStateKey maxStateKey = {.key = {.groupId = UINT64_MAX, .ts = INT64_MAX}, .opNum = INT64_MAX};
  STREAM_STATE_PUT_ROCKSDB(pState, "state", &maxStateKey, "", 0);
Y
yihaoDeng 已提交
1364 1365
  char             buf[128] = {0};
  int32_t          klen = stateKeyEncode((void*)&maxStateKey, buf);
dengyihao's avatar
dengyihao 已提交
1366
  SStreamStateCur* pCur = taosMemoryCalloc(1, sizeof(SStreamStateCur));
dengyihao's avatar
dengyihao 已提交
1367
  if (pCur == NULL) return NULL;
Y
yihaoDeng 已提交
1368
  pCur->db = wrapper->rocksdb;
1369 1370
  pCur->iter = streamStateIterCreate(pState, "state", (rocksdb_snapshot_t**)&pCur->snapshot,
                                     (rocksdb_readoptions_t**)&pCur->readOpt);
dengyihao's avatar
dengyihao 已提交
1371
  rocksdb_iter_seek(pCur->iter, buf, (size_t)klen);
dengyihao's avatar
dengyihao 已提交
1372

dengyihao's avatar
dengyihao 已提交
1373
  rocksdb_iter_prev(pCur->iter);
dengyihao's avatar
dengyihao 已提交
1374
  while (rocksdb_iter_valid(pCur->iter) && iterValueIsStale(pCur->iter)) {
dengyihao's avatar
dengyihao 已提交
1375 1376
    rocksdb_iter_prev(pCur->iter);
  }
dengyihao's avatar
dengyihao 已提交
1377

dengyihao's avatar
dengyihao 已提交
1378 1379 1380 1381
  if (!rocksdb_iter_valid(pCur->iter)) {
    streamStateFreeCur(pCur);
    pCur = NULL;
  }
dengyihao's avatar
dengyihao 已提交
1382
  STREAM_STATE_DEL_ROCKSDB(pState, "state", &maxStateKey);
dengyihao's avatar
dengyihao 已提交
1383
  return pCur;
dengyihao's avatar
dengyihao 已提交
1384 1385
}

dengyihao's avatar
dengyihao 已提交
1386
SStreamStateCur* streamStateGetCur_rocksdb(SStreamState* pState, const SWinKey* key) {
dengyihao's avatar
dengyihao 已提交
1387
  qDebug("streamStateGetCur_rocksdb");
Y
yihaoDeng 已提交
1388 1389
  SBackendCfWrapper* wrapper = pState->pTdbState->pBackendCfWrapper;
  SStreamStateCur*   pCur = taosMemoryCalloc(1, sizeof(SStreamStateCur));
dengyihao's avatar
dengyihao 已提交
1390

dengyihao's avatar
dengyihao 已提交
1391
  if (pCur == NULL) return NULL;
Y
yihaoDeng 已提交
1392
  pCur->db = wrapper->rocksdb;
1393 1394
  pCur->iter = streamStateIterCreate(pState, "state", (rocksdb_snapshot_t**)&pCur->snapshot,
                                     (rocksdb_readoptions_t**)&pCur->readOpt);
dengyihao's avatar
dengyihao 已提交
1395

dengyihao's avatar
dengyihao 已提交
1396
  SStateKey sKey = {.key = *key, .opNum = pState->number};
dengyihao's avatar
dengyihao 已提交
1397
  char      buf[128] = {0};
dengyihao's avatar
dengyihao 已提交
1398
  int       len = stateKeyEncode((void*)&sKey, buf);
dengyihao's avatar
dengyihao 已提交
1399

dengyihao's avatar
dengyihao 已提交
1400
  rocksdb_iter_seek(pCur->iter, buf, len);
dengyihao's avatar
dengyihao 已提交
1401

dengyihao's avatar
dengyihao 已提交
1402
  if (rocksdb_iter_valid(pCur->iter) && !iterValueIsStale(pCur->iter)) {
dengyihao's avatar
dengyihao 已提交
1403 1404
    size_t vlen;
    char*  val = (char*)rocksdb_iter_value(pCur->iter, &vlen);
dengyihao's avatar
dengyihao 已提交
1405 1406 1407 1408 1409 1410 1411 1412 1413 1414
    if (!streamStateValueIsStale(val)) {
      SStateKey curKey;
      size_t    kLen = 0;
      char*     keyStr = (char*)rocksdb_iter_key(pCur->iter, &kLen);
      stateKeyDecode((void*)&curKey, keyStr);

      if (stateKeyCmpr(&sKey, sizeof(sKey), &curKey, sizeof(curKey)) == 0) {
        pCur->number = pState->number;
        return pCur;
      }
dengyihao's avatar
dengyihao 已提交
1415
    }
dengyihao's avatar
dengyihao 已提交
1416
  }
dengyihao's avatar
dengyihao 已提交
1417 1418 1419
  streamStateFreeCur(pCur);
  return NULL;
}
dengyihao's avatar
dengyihao 已提交
1420

dengyihao's avatar
dengyihao 已提交
1421 1422 1423 1424 1425 1426 1427 1428 1429 1430 1431 1432 1433 1434
// func cf
int32_t streamStateFuncPut_rocksdb(SStreamState* pState, const STupleKey* key, const void* value, int32_t vLen) {
  int code = 0;
  STREAM_STATE_PUT_ROCKSDB(pState, "func", key, (void*)value, vLen);
  return code;
}
int32_t streamStateFuncGet_rocksdb(SStreamState* pState, const STupleKey* key, void** pVal, int32_t* pVLen) {
  int code = 0;
  STREAM_STATE_GET_ROCKSDB(pState, "func", key, pVal, pVLen);
  return 0;
}
int32_t streamStateFuncDel_rocksdb(SStreamState* pState, const STupleKey* key) {
  int code = 0;
  STREAM_STATE_DEL_ROCKSDB(pState, "func", key);
dengyihao's avatar
dengyihao 已提交
1435 1436
  return 0;
}
dengyihao's avatar
dengyihao 已提交
1437

dengyihao's avatar
dengyihao 已提交
1438
// session cf
dengyihao's avatar
dengyihao 已提交
1439 1440 1441
int32_t streamStateSessionPut_rocksdb(SStreamState* pState, const SSessionKey* key, const void* value, int32_t vLen) {
  int              code = 0;
  SStateSessionKey sKey = {.key = *key, .opNum = pState->number};
dengyihao's avatar
dengyihao 已提交
1442
  STREAM_STATE_PUT_ROCKSDB(pState, "sess", &sKey, value, vLen);
dengyihao's avatar
dengyihao 已提交
1443 1444
  if (code == -1) {
  }
dengyihao's avatar
dengyihao 已提交
1445 1446
  return code;
}
dengyihao's avatar
dengyihao 已提交
1447 1448 1449 1450 1451 1452 1453 1454 1455 1456
int32_t streamStateSessionGet_rocksdb(SStreamState* pState, SSessionKey* key, void** pVal, int32_t* pVLen) {
  qDebug("streamStateSessionGet_rocksdb");
  int              code = 0;
  SStreamStateCur* pCur = streamStateSessionSeekKeyCurrentNext_rocksdb(pState, key);
  SSessionKey      resKey = *key;
  void*            tmp = NULL;
  int32_t          vLen = 0;
  code = streamStateSessionGetKVByCur_rocksdb(pCur, &resKey, &tmp, &vLen);
  if (code == 0) {
    if (pVLen != NULL) *pVLen = vLen;
dengyihao's avatar
dengyihao 已提交
1457

dengyihao's avatar
dengyihao 已提交
1458 1459 1460 1461 1462 1463
    if (key->win.skey != resKey.win.skey) {
      code = -1;
    } else {
      *key = resKey;
      *pVal = taosMemoryCalloc(1, *pVLen);
      memcpy(*pVal, tmp, *pVLen);
dengyihao's avatar
dengyihao 已提交
1464 1465
    }
  }
dengyihao's avatar
dengyihao 已提交
1466
  taosMemoryFree(tmp);
dengyihao's avatar
dengyihao 已提交
1467 1468 1469
  streamStateFreeCur(pCur);
  // impl later
  return code;
dengyihao's avatar
dengyihao 已提交
1470
}
dengyihao's avatar
dengyihao 已提交
1471 1472 1473 1474 1475 1476 1477

int32_t streamStateSessionDel_rocksdb(SStreamState* pState, const SSessionKey* key) {
  int              code = 0;
  SStateSessionKey sKey = {.key = *key, .opNum = pState->number};
  STREAM_STATE_DEL_ROCKSDB(pState, "sess", &sKey);
  return code;
}
dengyihao's avatar
dengyihao 已提交
1478
SStreamStateCur* streamStateSessionSeekKeyCurrentPrev_rocksdb(SStreamState* pState, const SSessionKey* key) {
dengyihao's avatar
dengyihao 已提交
1479
  qDebug("streamStateSessionSeekKeyCurrentPrev_rocksdb");
Y
yihaoDeng 已提交
1480

Y
yihaoDeng 已提交
1481 1482
  SBackendCfWrapper* wrapper = pState->pTdbState->pBackendCfWrapper;
  SStreamStateCur*   pCur = taosMemoryCalloc(1, sizeof(SStreamStateCur));
dengyihao's avatar
dengyihao 已提交
1483 1484 1485 1486
  if (pCur == NULL) {
    return NULL;
  }
  pCur->number = pState->number;
Y
yihaoDeng 已提交
1487
  pCur->db = wrapper->rocksdb;
1488 1489
  pCur->iter = streamStateIterCreate(pState, "sess", (rocksdb_snapshot_t**)&pCur->snapshot,
                                     (rocksdb_readoptions_t**)&pCur->readOpt);
dengyihao's avatar
dengyihao 已提交
1490

dengyihao's avatar
dengyihao 已提交
1491
  char             buf[128] = {0};
dengyihao's avatar
dengyihao 已提交
1492 1493
  SStateSessionKey sKey = {.key = *key, .opNum = pState->number};
  int              len = stateSessionKeyEncode(&sKey, buf);
dengyihao's avatar
dengyihao 已提交
1494
  if (!streamStateIterSeekAndValid(pCur->iter, buf, len)) {
dengyihao's avatar
dengyihao 已提交
1495
    streamStateFreeCur(pCur);
dengyihao's avatar
dengyihao 已提交
1496
    return NULL;
dengyihao's avatar
dengyihao 已提交
1497
  }
dengyihao's avatar
dengyihao 已提交
1498 1499
  while (rocksdb_iter_valid(pCur->iter) && iterValueIsStale(pCur->iter)) rocksdb_iter_prev(pCur->iter);

dengyihao's avatar
dengyihao 已提交
1500
  if (!rocksdb_iter_valid(pCur->iter)) {
dengyihao's avatar
dengyihao 已提交
1501 1502
    streamStateFreeCur(pCur);
    return NULL;
dengyihao's avatar
dengyihao 已提交
1503 1504 1505 1506 1507 1508 1509 1510 1511 1512
  }

  int32_t          c = 0;
  size_t           klen;
  const char*      iKey = rocksdb_iter_key(pCur->iter, &klen);
  SStateSessionKey curKey = {0};
  stateSessionKeyDecode(&curKey, (char*)iKey);
  if (stateSessionKeyCmpr(&sKey, sizeof(sKey), &curKey, sizeof(curKey)) >= 0) return pCur;

  rocksdb_iter_prev(pCur->iter);
dengyihao's avatar
dengyihao 已提交
1513
  if (!rocksdb_iter_valid(pCur->iter)) {
dengyihao's avatar
dengyihao 已提交
1514 1515
    // qWarn("streamState failed to seek key prev
    // %s", toString);
dengyihao's avatar
dengyihao 已提交
1516 1517 1518
    streamStateFreeCur(pCur);
    return NULL;
  }
dengyihao's avatar
dengyihao 已提交
1519 1520
  return pCur;
}
dengyihao's avatar
dengyihao 已提交
1521
SStreamStateCur* streamStateSessionSeekKeyCurrentNext_rocksdb(SStreamState* pState, SSessionKey* key) {
dengyihao's avatar
dengyihao 已提交
1522
  qDebug("streamStateSessionSeekKeyCurrentNext_rocksdb");
Y
yihaoDeng 已提交
1523 1524
  SBackendCfWrapper* wrapper = pState->pTdbState->pBackendCfWrapper;
  SStreamStateCur*   pCur = taosMemoryCalloc(1, sizeof(SStreamStateCur));
dengyihao's avatar
dengyihao 已提交
1525 1526 1527
  if (pCur == NULL) {
    return NULL;
  }
Y
yihaoDeng 已提交
1528
  pCur->db = wrapper->rocksdb;
1529 1530
  pCur->iter = streamStateIterCreate(pState, "sess", (rocksdb_snapshot_t**)&pCur->snapshot,
                                     (rocksdb_readoptions_t**)&pCur->readOpt);
dengyihao's avatar
dengyihao 已提交
1531
  pCur->number = pState->number;
dengyihao's avatar
dengyihao 已提交
1532

dengyihao's avatar
dengyihao 已提交
1533
  char buf[128] = {0};
dengyihao's avatar
dengyihao 已提交
1534

dengyihao's avatar
dengyihao 已提交
1535
  SStateSessionKey sKey = {.key = *key, .opNum = pState->number};
dengyihao's avatar
dengyihao 已提交
1536
  int              len = stateSessionKeyEncode(&sKey, buf);
dengyihao's avatar
dengyihao 已提交
1537
  if (!streamStateIterSeekAndValid(pCur->iter, buf, len)) {
dengyihao's avatar
dengyihao 已提交
1538 1539
    streamStateFreeCur(pCur);
    return NULL;
dengyihao's avatar
dengyihao 已提交
1540
  }
dengyihao's avatar
dengyihao 已提交
1541
  if (iterValueIsStale(pCur->iter)) {
dengyihao's avatar
dengyihao 已提交
1542 1543 1544
    streamStateFreeCur(pCur);
    return NULL;
  }
dengyihao's avatar
dengyihao 已提交
1545 1546 1547 1548 1549 1550 1551
  size_t           klen;
  const char*      iKey = rocksdb_iter_key(pCur->iter, &klen);
  SStateSessionKey curKey = {0};
  stateSessionKeyDecode(&curKey, (char*)iKey);
  if (stateSessionKeyCmpr(&sKey, sizeof(sKey), &curKey, sizeof(curKey)) <= 0) return pCur;

  rocksdb_iter_next(pCur->iter);
dengyihao's avatar
dengyihao 已提交
1552 1553 1554 1555
  if (!rocksdb_iter_valid(pCur->iter)) {
    streamStateFreeCur(pCur);
    return NULL;
  }
dengyihao's avatar
dengyihao 已提交
1556 1557
  return pCur;
}
dengyihao's avatar
dengyihao 已提交
1558

dengyihao's avatar
dengyihao 已提交
1559
SStreamStateCur* streamStateSessionSeekKeyNext_rocksdb(SStreamState* pState, const SSessionKey* key) {
dengyihao's avatar
dengyihao 已提交
1560
  qDebug("streamStateSessionSeekKeyNext_rocksdb");
Y
yihaoDeng 已提交
1561 1562
  SBackendCfWrapper* wrapper = pState->pTdbState->pBackendCfWrapper;
  SStreamStateCur*   pCur = taosMemoryCalloc(1, sizeof(SStreamStateCur));
dengyihao's avatar
dengyihao 已提交
1563 1564 1565
  if (pCur == NULL) {
    return NULL;
  }
Y
yihaoDeng 已提交
1566
  pCur->db = wrapper->rocksdb;
1567 1568
  pCur->iter = streamStateIterCreate(pState, "sess", (rocksdb_snapshot_t**)&pCur->snapshot,
                                     (rocksdb_readoptions_t**)&pCur->readOpt);
dengyihao's avatar
dengyihao 已提交
1569
  pCur->number = pState->number;
dengyihao's avatar
dengyihao 已提交
1570

dengyihao's avatar
dengyihao 已提交
1571
  SStateSessionKey sKey = {.key = *key, .opNum = pState->number};
dengyihao's avatar
dengyihao 已提交
1572

dengyihao's avatar
dengyihao 已提交
1573
  char buf[128] = {0};
dengyihao's avatar
dengyihao 已提交
1574
  int  len = stateSessionKeyEncode(&sKey, buf);
dengyihao's avatar
dengyihao 已提交
1575
  if (!streamStateIterSeekAndValid(pCur->iter, buf, len)) {
dengyihao's avatar
dengyihao 已提交
1576 1577
    streamStateFreeCur(pCur);
    return NULL;
dengyihao's avatar
dengyihao 已提交
1578
  }
dengyihao's avatar
dengyihao 已提交
1579 1580 1581 1582 1583
  while (rocksdb_iter_valid(pCur->iter) && iterValueIsStale(pCur->iter)) rocksdb_iter_next(pCur->iter);
  if (!rocksdb_iter_valid(pCur->iter)) {
    streamStateFreeCur(pCur);
    return NULL;
  }
dengyihao's avatar
dengyihao 已提交
1584 1585 1586 1587 1588 1589 1590
  size_t           klen;
  const char*      iKey = rocksdb_iter_key(pCur->iter, &klen);
  SStateSessionKey curKey = {0};
  stateSessionKeyDecode(&curKey, (char*)iKey);
  if (stateSessionKeyCmpr(&sKey, sizeof(sKey), &curKey, sizeof(curKey)) < 0) return pCur;

  rocksdb_iter_next(pCur->iter);
dengyihao's avatar
dengyihao 已提交
1591 1592 1593 1594
  if (!rocksdb_iter_valid(pCur->iter)) {
    streamStateFreeCur(pCur);
    return NULL;
  }
dengyihao's avatar
dengyihao 已提交
1595 1596
  return pCur;
}
dengyihao's avatar
dengyihao 已提交
1597
int32_t streamStateSessionGetKVByCur_rocksdb(SStreamStateCur* pCur, SSessionKey* pKey, void** pVal, int32_t* pVLen) {
dengyihao's avatar
dengyihao 已提交
1598
  qDebug("streamStateSessionGetKVByCur_rocksdb");
dengyihao's avatar
dengyihao 已提交
1599 1600 1601
  if (!pCur) {
    return -1;
  }
dengyihao's avatar
dengyihao 已提交
1602
  SStateSessionKey ktmp = {0};
dengyihao's avatar
dengyihao 已提交
1603
  size_t           kLen = 0, vLen = 0;
dengyihao's avatar
dengyihao 已提交
1604

dengyihao's avatar
dengyihao 已提交
1605
  if (!rocksdb_iter_valid(pCur->iter) || iterValueIsStale(pCur->iter)) {
dengyihao's avatar
dengyihao 已提交
1606 1607 1608 1609 1610
    return -1;
  }
  const char* curKey = rocksdb_iter_key(pCur->iter, (size_t*)&kLen);
  stateSessionKeyDecode((void*)&ktmp, (char*)curKey);

dengyihao's avatar
dengyihao 已提交
1611
  SStateSessionKey* pKTmp = &ktmp;
dengyihao's avatar
dengyihao 已提交
1612 1613 1614 1615 1616
  const char*       vval = rocksdb_iter_value(pCur->iter, (size_t*)&vLen);
  char*             val = NULL;
  int32_t           len = decodeValueFunc((void*)vval, vLen, NULL, &val);
  if (len < 0) {
    return -1;
dengyihao's avatar
dengyihao 已提交
1617
  }
dengyihao's avatar
dengyihao 已提交
1618 1619 1620 1621 1622
  if (pVal != NULL) {
    *pVal = (char*)val;
  } else {
    taosMemoryFree(val);
  }
dengyihao's avatar
dengyihao 已提交
1623
  if (pVLen != NULL) *pVLen = len;
dengyihao's avatar
dengyihao 已提交
1624 1625 1626 1627 1628 1629 1630 1631 1632 1633

  if (pKTmp->opNum != pCur->number) {
    return -1;
  }
  if (pKey->groupId != 0 && pKey->groupId != pKTmp->key.groupId) {
    return -1;
  }
  *pKey = pKTmp->key;
  return 0;
}
dengyihao's avatar
dengyihao 已提交
1634 1635 1636
// fill cf
int32_t streamStateFillPut_rocksdb(SStreamState* pState, const SWinKey* key, const void* value, int32_t vLen) {
  int code = 0;
dengyihao's avatar
dengyihao 已提交
1637

dengyihao's avatar
dengyihao 已提交
1638 1639 1640 1641 1642 1643 1644 1645 1646 1647 1648 1649 1650 1651
  STREAM_STATE_PUT_ROCKSDB(pState, "fill", key, value, vLen);
  return code;
}

int32_t streamStateFillGet_rocksdb(SStreamState* pState, const SWinKey* key, void** pVal, int32_t* pVLen) {
  int code = 0;
  STREAM_STATE_GET_ROCKSDB(pState, "fill", key, pVal, pVLen);
  return code;
}
int32_t streamStateFillDel_rocksdb(SStreamState* pState, const SWinKey* key) {
  int code = 0;
  STREAM_STATE_DEL_ROCKSDB(pState, "fill", key);
  return code;
}
dengyihao's avatar
dengyihao 已提交
1652

dengyihao's avatar
dengyihao 已提交
1653 1654
SStreamStateCur* streamStateFillGetCur_rocksdb(SStreamState* pState, const SWinKey* key) {
  qDebug("streamStateFillGetCur_rocksdb");
Y
yihaoDeng 已提交
1655 1656
  SStreamStateCur*   pCur = taosMemoryCalloc(1, sizeof(SStreamStateCur));
  SBackendCfWrapper* wrapper = pState->pTdbState->pBackendCfWrapper;
dengyihao's avatar
dengyihao 已提交
1657 1658 1659

  if (pCur == NULL) return NULL;

Y
yihaoDeng 已提交
1660
  pCur->db = wrapper->rocksdb;
1661 1662
  pCur->iter = streamStateIterCreate(pState, "fill", (rocksdb_snapshot_t**)&pCur->snapshot,
                                     (rocksdb_readoptions_t**)&pCur->readOpt);
dengyihao's avatar
dengyihao 已提交
1663

dengyihao's avatar
dengyihao 已提交
1664 1665
  char buf[128] = {0};
  int  len = winKeyEncode((void*)key, buf);
dengyihao's avatar
dengyihao 已提交
1666 1667 1668
  if (!streamStateIterSeekAndValid(pCur->iter, buf, len)) {
    streamStateFreeCur(pCur);
    return NULL;
dengyihao's avatar
dengyihao 已提交
1669
  }
dengyihao's avatar
dengyihao 已提交
1670 1671 1672
  if (iterValueIsStale(pCur->iter)) {
    streamStateFreeCur(pCur);
    return NULL;
dengyihao's avatar
dengyihao 已提交
1673
  }
dengyihao's avatar
dengyihao 已提交
1674

dengyihao's avatar
dengyihao 已提交
1675
  if (rocksdb_iter_valid(pCur->iter)) {
dengyihao's avatar
dengyihao 已提交
1676 1677 1678 1679 1680
    size_t  kLen;
    SWinKey curKey;
    char*   keyStr = (char*)rocksdb_iter_key(pCur->iter, &kLen);
    winKeyDecode((void*)&curKey, keyStr);
    if (winKeyCmpr(key, sizeof(*key), &curKey, sizeof(curKey)) == 0) {
dengyihao's avatar
dengyihao 已提交
1681 1682 1683
      return pCur;
    }
  }
dengyihao's avatar
dengyihao 已提交
1684

dengyihao's avatar
dengyihao 已提交
1685 1686 1687
  streamStateFreeCur(pCur);
  return NULL;
}
dengyihao's avatar
dengyihao 已提交
1688 1689 1690 1691 1692 1693
int32_t streamStateFillGetKVByCur_rocksdb(SStreamStateCur* pCur, SWinKey* pKey, const void** pVal, int32_t* pVLen) {
  qDebug("streamStateFillGetKVByCur_rocksdb");
  if (!pCur) {
    return -1;
  }
  SWinKey winKey;
dengyihao's avatar
dengyihao 已提交
1694 1695 1696
  if (!rocksdb_iter_valid(pCur->iter) || iterValueIsStale(pCur->iter)) {
    return -1;
  }
dengyihao's avatar
dengyihao 已提交
1697 1698
  size_t klen, vlen;
  char*  keyStr = (char*)rocksdb_iter_key(pCur->iter, &klen);
dengyihao's avatar
dengyihao 已提交
1699
  winKeyDecode(&winKey, keyStr);
dengyihao's avatar
dengyihao 已提交
1700

dengyihao's avatar
dengyihao 已提交
1701
  const char* valStr = rocksdb_iter_value(pCur->iter, &vlen);
dengyihao's avatar
dengyihao 已提交
1702 1703
  // char*       dst = NULL;
  int32_t len = decodeValueFunc((void*)valStr, vlen, NULL, (char**)pVal);
dengyihao's avatar
dengyihao 已提交
1704
  if (len < 0) {
dengyihao's avatar
dengyihao 已提交
1705 1706
    return -1;
  }
dengyihao's avatar
dengyihao 已提交
1707
  if (pVLen != NULL) *pVLen = len;
dengyihao's avatar
dengyihao 已提交
1708

dengyihao's avatar
dengyihao 已提交
1709 1710 1711 1712
  *pKey = winKey;
  return 0;
}

dengyihao's avatar
dengyihao 已提交
1713
SStreamStateCur* streamStateFillSeekKeyNext_rocksdb(SStreamState* pState, const SWinKey* key) {
dengyihao's avatar
dengyihao 已提交
1714
  qDebug("streamStateFillSeekKeyNext_rocksdb");
Y
yihaoDeng 已提交
1715 1716
  SBackendCfWrapper* wrapper = pState->pTdbState->pBackendCfWrapper;
  SStreamStateCur*   pCur = taosMemoryCalloc(1, sizeof(SStreamStateCur));
dengyihao's avatar
dengyihao 已提交
1717 1718 1719
  if (!pCur) {
    return NULL;
  }
dengyihao's avatar
dengyihao 已提交
1720

Y
yihaoDeng 已提交
1721
  pCur->db = wrapper->rocksdb;
1722 1723
  pCur->iter = streamStateIterCreate(pState, "fill", (rocksdb_snapshot_t**)&pCur->snapshot,
                                     (rocksdb_readoptions_t**)&pCur->readOpt);
dengyihao's avatar
dengyihao 已提交
1724

dengyihao's avatar
dengyihao 已提交
1725
  char buf[128] = {0};
dengyihao's avatar
dengyihao 已提交
1726
  int  len = winKeyEncode((void*)key, buf);
dengyihao's avatar
dengyihao 已提交
1727 1728 1729
  if (!streamStateIterSeekAndValid(pCur->iter, buf, len)) {
    streamStateFreeCur(pCur);
    return NULL;
dengyihao's avatar
dengyihao 已提交
1730
  }
dengyihao's avatar
dengyihao 已提交
1731 1732 1733 1734
  // skip stale data
  while (rocksdb_iter_valid(pCur->iter) && iterValueIsStale(pCur->iter)) {
    rocksdb_iter_next(pCur->iter);
  }
dengyihao's avatar
dengyihao 已提交
1735

dengyihao's avatar
dengyihao 已提交
1736
  if (rocksdb_iter_valid(pCur->iter)) {
dengyihao's avatar
dengyihao 已提交
1737 1738 1739 1740 1741 1742 1743 1744
    SWinKey curKey;
    size_t  kLen = 0;
    char*   keyStr = (char*)rocksdb_iter_key(pCur->iter, &kLen);
    winKeyDecode((void*)&curKey, keyStr);
    if (winKeyCmpr(key, sizeof(*key), &curKey, sizeof(curKey)) > 0) {
      return pCur;
    }
    rocksdb_iter_next(pCur->iter);
dengyihao's avatar
dengyihao 已提交
1745
    return pCur;
dengyihao's avatar
dengyihao 已提交
1746 1747 1748 1749 1750
  }
  streamStateFreeCur(pCur);
  return NULL;
}
SStreamStateCur* streamStateFillSeekKeyPrev_rocksdb(SStreamState* pState, const SWinKey* key) {
dengyihao's avatar
dengyihao 已提交
1751
  qDebug("streamStateFillSeekKeyPrev_rocksdb");
Y
yihaoDeng 已提交
1752 1753
  SBackendCfWrapper* wrapper = pState->pTdbState->pBackendCfWrapper;
  SStreamStateCur*   pCur = taosMemoryCalloc(1, sizeof(SStreamStateCur));
dengyihao's avatar
dengyihao 已提交
1754 1755 1756
  if (pCur == NULL) {
    return NULL;
  }
dengyihao's avatar
dengyihao 已提交
1757

Y
yihaoDeng 已提交
1758
  pCur->db = wrapper->rocksdb;
1759 1760
  pCur->iter = streamStateIterCreate(pState, "fill", (rocksdb_snapshot_t**)&pCur->snapshot,
                                     (rocksdb_readoptions_t**)&pCur->readOpt);
dengyihao's avatar
dengyihao 已提交
1761

dengyihao's avatar
dengyihao 已提交
1762
  char buf[128] = {0};
dengyihao's avatar
dengyihao 已提交
1763
  int  len = winKeyEncode((void*)key, buf);
dengyihao's avatar
dengyihao 已提交
1764 1765 1766
  if (!streamStateIterSeekAndValid(pCur->iter, buf, len)) {
    streamStateFreeCur(pCur);
    return NULL;
dengyihao's avatar
dengyihao 已提交
1767
  }
dengyihao's avatar
dengyihao 已提交
1768 1769 1770
  while (rocksdb_iter_valid(pCur->iter) && iterValueIsStale(pCur->iter)) {
    rocksdb_iter_prev(pCur->iter);
  }
dengyihao's avatar
dengyihao 已提交
1771

dengyihao's avatar
dengyihao 已提交
1772
  if (rocksdb_iter_valid(pCur->iter)) {
dengyihao's avatar
dengyihao 已提交
1773 1774 1775 1776 1777 1778 1779 1780
    SWinKey curKey;
    size_t  kLen = 0;
    char*   keyStr = (char*)rocksdb_iter_key(pCur->iter, &kLen);
    winKeyDecode((void*)&curKey, keyStr);
    if (winKeyCmpr(key, sizeof(*key), &curKey, sizeof(curKey)) < 0) {
      return pCur;
    }
    rocksdb_iter_prev(pCur->iter);
dengyihao's avatar
dengyihao 已提交
1781
    return pCur;
dengyihao's avatar
dengyihao 已提交
1782
  }
dengyihao's avatar
dengyihao 已提交
1783

dengyihao's avatar
dengyihao 已提交
1784 1785 1786
  streamStateFreeCur(pCur);
  return NULL;
}
dengyihao's avatar
dengyihao 已提交
1787
int32_t streamStateSessionGetKeyByRange_rocksdb(SStreamState* pState, const SSessionKey* key, SSessionKey* curKey) {
dengyihao's avatar
dengyihao 已提交
1788
  qDebug("streamStateSessionGetKeyByRange_rocksdb");
Y
yihaoDeng 已提交
1789 1790
  SBackendCfWrapper* wrapper = pState->pTdbState->pBackendCfWrapper;
  SStreamStateCur*   pCur = taosMemoryCalloc(1, sizeof(SStreamStateCur));
dengyihao's avatar
dengyihao 已提交
1791 1792 1793 1794
  if (pCur == NULL) {
    return -1;
  }
  pCur->number = pState->number;
Y
yihaoDeng 已提交
1795
  pCur->db = wrapper->rocksdb;
1796 1797
  pCur->iter = streamStateIterCreate(pState, "sess", (rocksdb_snapshot_t**)&pCur->snapshot,
                                     (rocksdb_readoptions_t**)&pCur->readOpt);
dengyihao's avatar
dengyihao 已提交
1798 1799 1800

  SStateSessionKey sKey = {.key = *key, .opNum = pState->number};
  int32_t          c = 0;
dengyihao's avatar
dengyihao 已提交
1801
  char             buf[128] = {0};
dengyihao's avatar
dengyihao 已提交
1802
  int              len = stateSessionKeyEncode(&sKey, buf);
dengyihao's avatar
dengyihao 已提交
1803 1804 1805
  if (!streamStateIterSeekAndValid(pCur->iter, buf, len)) {
    streamStateFreeCur(pCur);
    return -1;
dengyihao's avatar
dengyihao 已提交
1806 1807
  }

dengyihao's avatar
dengyihao 已提交
1808
  size_t           kLen;
dengyihao's avatar
dengyihao 已提交
1809 1810 1811 1812 1813 1814 1815 1816 1817 1818 1819 1820 1821 1822 1823 1824 1825 1826 1827 1828 1829 1830 1831 1832 1833 1834 1835 1836 1837 1838 1839 1840 1841 1842 1843 1844 1845 1846
  const char*      iKeyStr = rocksdb_iter_key(pCur->iter, (size_t*)&kLen);
  SStateSessionKey iKey = {0};
  stateSessionKeyDecode(&iKey, (char*)iKeyStr);

  c = stateSessionKeyCmpr(&sKey, sizeof(sKey), &iKey, sizeof(iKey));

  SSessionKey resKey = *key;
  int32_t     code = streamStateSessionGetKVByCur_rocksdb(pCur, &resKey, NULL, 0);
  if (code == 0 && sessionRangeKeyCmpr(key, &resKey) == 0) {
    *curKey = resKey;
    streamStateFreeCur(pCur);
    return code;
  }

  if (c > 0) {
    streamStateCurNext_rocksdb(pState, pCur);
    code = streamStateSessionGetKVByCur_rocksdb(pCur, &resKey, NULL, 0);
    if (code == 0 && sessionRangeKeyCmpr(key, &resKey) == 0) {
      *curKey = resKey;
      streamStateFreeCur(pCur);
      return code;
    }
  } else if (c < 0) {
    streamStateCurPrev(pState, pCur);
    code = streamStateSessionGetKVByCur_rocksdb(pCur, &resKey, NULL, 0);
    if (code == 0 && sessionRangeKeyCmpr(key, &resKey) == 0) {
      *curKey = resKey;
      streamStateFreeCur(pCur);
      return code;
    }
  }

  streamStateFreeCur(pCur);
  return -1;
}

int32_t streamStateSessionAddIfNotExist_rocksdb(SStreamState* pState, SSessionKey* key, TSKEY gap, void** pVal,
                                                int32_t* pVLen) {
dengyihao's avatar
dengyihao 已提交
1847
  qDebug("streamStateSessionAddIfNotExist_rocksdb");
dengyihao's avatar
dengyihao 已提交
1848 1849 1850 1851 1852 1853 1854 1855 1856 1857 1858
  // todo refactor
  int32_t     res = 0;
  SSessionKey originKey = *key;
  SSessionKey searchKey = *key;
  searchKey.win.skey = key->win.skey - gap;
  searchKey.win.ekey = key->win.ekey + gap;
  int32_t valSize = *pVLen;

  void* tmp = taosMemoryMalloc(valSize);

  SStreamStateCur* pCur = streamStateSessionSeekKeyCurrentPrev_rocksdb(pState, key);
dengyihao's avatar
dengyihao 已提交
1859 1860 1861 1862
  if (pCur == NULL) {
  }
  int32_t code = streamStateSessionGetKVByCur_rocksdb(pCur, key, pVal, pVLen);

dengyihao's avatar
dengyihao 已提交
1863 1864 1865
  if (code == 0) {
    if (sessionRangeKeyCmpr(&searchKey, key) == 0) {
      memcpy(tmp, *pVal, valSize);
dengyihao's avatar
dengyihao 已提交
1866
      taosMemoryFreeClear(*pVal);
dengyihao's avatar
dengyihao 已提交
1867 1868 1869
      streamStateSessionDel_rocksdb(pState, key);
      goto _end;
    }
dengyihao's avatar
dengyihao 已提交
1870
    taosMemoryFreeClear(*pVal);
dengyihao's avatar
dengyihao 已提交
1871 1872 1873 1874
    streamStateCurNext_rocksdb(pState, pCur);
  } else {
    *key = originKey;
    streamStateFreeCur(pCur);
dengyihao's avatar
dengyihao 已提交
1875
    taosMemoryFreeClear(*pVal);
dengyihao's avatar
dengyihao 已提交
1876 1877 1878 1879 1880 1881 1882 1883 1884 1885 1886 1887 1888 1889 1890 1891 1892
    pCur = streamStateSessionSeekKeyNext_rocksdb(pState, key);
  }

  code = streamStateSessionGetKVByCur_rocksdb(pCur, key, pVal, pVLen);
  if (code == 0) {
    if (sessionRangeKeyCmpr(&searchKey, key) == 0) {
      memcpy(tmp, *pVal, valSize);
      streamStateSessionDel_rocksdb(pState, key);
      goto _end;
    }
  }

  *key = originKey;
  res = 1;
  memset(tmp, 0, valSize);

_end:
dengyihao's avatar
dengyihao 已提交
1893
  taosMemoryFree(*pVal);
dengyihao's avatar
dengyihao 已提交
1894 1895 1896 1897
  *pVal = tmp;
  streamStateFreeCur(pCur);
  return res;
}
dengyihao's avatar
dengyihao 已提交
1898 1899 1900 1901 1902 1903 1904 1905 1906 1907 1908 1909
int32_t streamStateSessionClear_rocksdb(SStreamState* pState) {
  qDebug("streamStateSessionClear_rocksdb");
  SSessionKey      key = {.win.skey = 0, .win.ekey = 0, .groupId = 0};
  SStreamStateCur* pCur = streamStateSessionSeekKeyCurrentNext_rocksdb(pState, &key);

  while (1) {
    SSessionKey delKey = {0};
    void*       buf = NULL;
    int32_t     size = 0;
    int32_t     code = streamStateSessionGetKVByCur_rocksdb(pCur, &delKey, &buf, &size);
    if (code == 0 && size > 0) {
      memset(buf, 0, size);
dengyihao's avatar
dengyihao 已提交
1910
      // refactor later
dengyihao's avatar
dengyihao 已提交
1911 1912
      streamStateSessionPut_rocksdb(pState, &delKey, buf, size);
    } else {
dengyihao's avatar
dengyihao 已提交
1913
      taosMemoryFreeClear(buf);
dengyihao's avatar
dengyihao 已提交
1914 1915
      break;
    }
dengyihao's avatar
dengyihao 已提交
1916 1917
    taosMemoryFreeClear(buf);

dengyihao's avatar
dengyihao 已提交
1918 1919 1920 1921 1922
    streamStateCurNext_rocksdb(pState, pCur);
  }
  streamStateFreeCur(pCur);
  return -1;
}
dengyihao's avatar
dengyihao 已提交
1923 1924
int32_t streamStateStateAddIfNotExist_rocksdb(SStreamState* pState, SSessionKey* key, char* pKeyData,
                                              int32_t keyDataLen, state_key_cmpr_fn fn, void** pVal, int32_t* pVLen) {
dengyihao's avatar
dengyihao 已提交
1925
  qDebug("streamStateStateAddIfNotExist_rocksdb");
dengyihao's avatar
dengyihao 已提交
1926 1927 1928 1929 1930 1931 1932 1933 1934 1935 1936 1937 1938 1939 1940
  // todo refactor
  int32_t     res = 0;
  SSessionKey tmpKey = *key;
  int32_t     valSize = *pVLen;
  void*       tmp = taosMemoryMalloc(valSize);
  // tdbRealloc(NULL, valSize);
  if (!tmp) {
    return -1;
  }

  SStreamStateCur* pCur = streamStateSessionSeekKeyCurrentPrev_rocksdb(pState, key);
  int32_t          code = streamStateSessionGetKVByCur_rocksdb(pCur, key, pVal, pVLen);
  if (code == 0) {
    if (key->win.skey <= tmpKey.win.skey && tmpKey.win.ekey <= key->win.ekey) {
      memcpy(tmp, *pVal, valSize);
dengyihao's avatar
dengyihao 已提交
1941
      streamStateSessionDel_rocksdb(pState, key);
dengyihao's avatar
dengyihao 已提交
1942 1943 1944 1945 1946 1947
      goto _end;
    }

    void* stateKey = (char*)(*pVal) + (valSize - keyDataLen);
    if (fn(pKeyData, stateKey) == true) {
      memcpy(tmp, *pVal, valSize);
dengyihao's avatar
dengyihao 已提交
1948
      streamStateSessionDel_rocksdb(pState, key);
dengyihao's avatar
dengyihao 已提交
1949 1950 1951 1952 1953 1954 1955 1956 1957
      goto _end;
    }

    streamStateCurNext_rocksdb(pState, pCur);
  } else {
    *key = tmpKey;
    streamStateFreeCur(pCur);
    pCur = streamStateSessionSeekKeyNext_rocksdb(pState, key);
  }
dengyihao's avatar
dengyihao 已提交
1958
  taosMemoryFreeClear(*pVal);
dengyihao's avatar
dengyihao 已提交
1959 1960 1961 1962 1963 1964 1965 1966 1967
  code = streamStateSessionGetKVByCur_rocksdb(pCur, key, pVal, pVLen);
  if (code == 0) {
    void* stateKey = (char*)(*pVal) + (valSize - keyDataLen);
    if (fn(pKeyData, stateKey) == true) {
      memcpy(tmp, *pVal, valSize);
      streamStateSessionDel_rocksdb(pState, key);
      goto _end;
    }
  }
dengyihao's avatar
dengyihao 已提交
1968
  taosMemoryFreeClear(*pVal);
dengyihao's avatar
dengyihao 已提交
1969 1970 1971 1972 1973 1974

  *key = tmpKey;
  res = 1;
  memset(tmp, 0, valSize);

_end:
dengyihao's avatar
dengyihao 已提交
1975
  taosMemoryFreeClear(*pVal);
dengyihao's avatar
dengyihao 已提交
1976 1977 1978 1979
  *pVal = tmp;
  streamStateFreeCur(pCur);
  return res;
}
dengyihao's avatar
dengyihao 已提交
1980

dengyihao's avatar
dengyihao 已提交
1981
//  partag cf
dengyihao's avatar
dengyihao 已提交
1982 1983 1984 1985 1986 1987 1988 1989 1990 1991 1992
int32_t streamStatePutParTag_rocksdb(SStreamState* pState, int64_t groupId, const void* tag, int32_t tagLen) {
  int code = 0;
  STREAM_STATE_PUT_ROCKSDB(pState, "partag", &groupId, tag, tagLen);
  return code;
}

int32_t streamStateGetParTag_rocksdb(SStreamState* pState, int64_t groupId, void** tagVal, int32_t* tagLen) {
  int code = 0;
  STREAM_STATE_GET_ROCKSDB(pState, "partag", &groupId, tagVal, tagLen);
  return code;
}
dengyihao's avatar
dengyihao 已提交
1993
// parname cfg
dengyihao's avatar
dengyihao 已提交
1994 1995
int32_t streamStatePutParName_rocksdb(SStreamState* pState, int64_t groupId, const char tbname[TSDB_TABLE_NAME_LEN]) {
  int code = 0;
dengyihao's avatar
dengyihao 已提交
1996
  STREAM_STATE_PUT_ROCKSDB(pState, "parname", &groupId, (char*)tbname, TSDB_TABLE_NAME_LEN);
dengyihao's avatar
dengyihao 已提交
1997 1998 1999 2000 2001 2002 2003 2004 2005
  return code;
}
int32_t streamStateGetParName_rocksdb(SStreamState* pState, int64_t groupId, void** pVal) {
  int    code = 0;
  size_t tagLen;
  STREAM_STATE_GET_ROCKSDB(pState, "parname", &groupId, pVal, &tagLen);
  return code;
}

dengyihao's avatar
dengyihao 已提交
2006 2007
int32_t streamDefaultPut_rocksdb(SStreamState* pState, const void* key, void* pVal, int32_t pVLen) {
  int code = 0;
dengyihao's avatar
dengyihao 已提交
2008
  STREAM_STATE_PUT_ROCKSDB(pState, "default", key, pVal, pVLen);
dengyihao's avatar
dengyihao 已提交
2009 2010 2011 2012
  return code;
}
int32_t streamDefaultGet_rocksdb(SStreamState* pState, const void* key, void** pVal, int32_t* pVLen) {
  int code = 0;
dengyihao's avatar
dengyihao 已提交
2013
  STREAM_STATE_GET_ROCKSDB(pState, "default", key, pVal, pVLen);
dengyihao's avatar
dengyihao 已提交
2014 2015 2016 2017
  return code;
}
int32_t streamDefaultDel_rocksdb(SStreamState* pState, const void* key) {
  int code = 0;
dengyihao's avatar
dengyihao 已提交
2018
  STREAM_STATE_DEL_ROCKSDB(pState, "default", key);
dengyihao's avatar
dengyihao 已提交
2019 2020 2021 2022 2023 2024 2025
  return code;
}

int32_t streamDefaultIterGet_rocksdb(SStreamState* pState, const void* start, const void* end, SArray* result) {
  int   code = 0;
  char* err = NULL;

Y
yihaoDeng 已提交
2026
  SBackendCfWrapper*     wrapper = pState->pTdbState->pBackendCfWrapper;
dengyihao's avatar
dengyihao 已提交
2027 2028 2029 2030 2031 2032 2033 2034 2035 2036 2037 2038 2039 2040 2041 2042 2043 2044 2045 2046 2047 2048 2049 2050 2051 2052 2053 2054 2055 2056 2057 2058
  rocksdb_snapshot_t*    snapshot = NULL;
  rocksdb_readoptions_t* readopts = NULL;
  rocksdb_iterator_t*    pIter = streamStateIterCreate(pState, "default", &snapshot, &readopts);
  if (pIter == NULL) {
    return -1;
  }

  rocksdb_iter_seek(pIter, start, strlen(start));
  while (rocksdb_iter_valid(pIter)) {
    const char* key = rocksdb_iter_key(pIter, NULL);
    int32_t     vlen = 0;
    const char* vval = rocksdb_iter_value(pIter, (size_t*)&vlen);
    char*       val = NULL;
    int32_t     len = decodeValueFunc((void*)vval, vlen, NULL, NULL);
    if (len < 0) {
      rocksdb_iter_next(pIter);
      continue;
    }

    if (end != NULL && strcmp(key, end) > 0) {
      break;
    }
    if (strncmp(key, start, strlen(start)) == 0 && strlen(key) >= strlen(start) + 1) {
      int64_t checkPoint = 0;
      if (sscanf(key + strlen(key), ":%" PRId64 "", &checkPoint) == 1) {
        taosArrayPush(result, &checkPoint);
      }
    } else {
      break;
    }
    rocksdb_iter_next(pIter);
  }
Y
yihaoDeng 已提交
2059
  rocksdb_release_snapshot(wrapper->rocksdb, snapshot);
dengyihao's avatar
dengyihao 已提交
2060 2061 2062 2063 2064
  rocksdb_readoptions_destroy(readopts);
  rocksdb_iter_destroy(pIter);
  return code;
}
void* streamDefaultIterCreate_rocksdb(SStreamState* pState) {
Y
yihaoDeng 已提交
2065 2066
  SStreamStateCur*   pCur = taosMemoryCalloc(1, sizeof(SStreamStateCur));
  SBackendCfWrapper* wrapper = pState->pTdbState->pBackendCfWrapper;
dengyihao's avatar
dengyihao 已提交
2067

Y
yihaoDeng 已提交
2068
  pCur->db = wrapper->rocksdb;
2069 2070
  pCur->iter = streamStateIterCreate(pState, "default", (rocksdb_snapshot_t**)&pCur->snapshot,
                                     (rocksdb_readoptions_t**)&pCur->readOpt);
dengyihao's avatar
dengyihao 已提交
2071 2072 2073 2074 2075 2076
  return pCur;
}
int32_t streamDefaultIterValid_rocksdb(void* iter) {
  SStreamStateCur* pCur = iter;
  bool             val = rocksdb_iter_valid(pCur->iter);

dengyihao's avatar
dengyihao 已提交
2077
  return val ? 1 : 0;
dengyihao's avatar
dengyihao 已提交
2078 2079 2080 2081 2082 2083 2084 2085 2086 2087 2088 2089 2090 2091 2092 2093 2094 2095 2096 2097 2098 2099 2100 2101 2102 2103 2104 2105 2106 2107 2108 2109 2110 2111 2112 2113
}
void streamDefaultIterSeek_rocksdb(void* iter, const char* key) {
  SStreamStateCur* pCur = iter;
  rocksdb_iter_seek(pCur->iter, key, strlen(key));
}
void streamDefaultIterNext_rocksdb(void* iter) {
  SStreamStateCur* pCur = iter;
  rocksdb_iter_next(pCur->iter);
}
char* streamDefaultIterKey_rocksdb(void* iter, int32_t* len) {
  SStreamStateCur* pCur = iter;
  return (char*)rocksdb_iter_key(pCur->iter, (size_t*)len);
}
char* streamDefaultIterVal_rocksdb(void* iter, int32_t* len) {
  SStreamStateCur* pCur = iter;
  int32_t          vlen = 0;
  char*            dst = NULL;
  const char*      vval = rocksdb_iter_value(pCur->iter, (size_t*)&vlen);
  if (decodeValueFunc((void*)vval, vlen, NULL, &dst) < 0) {
    return NULL;
  }
  return dst;
}
// batch func
void* streamStateCreateBatch() {
  rocksdb_writebatch_t* pBatch = rocksdb_writebatch_create();
  return pBatch;
}
int32_t streamStateGetBatchSize(void* pBatch) {
  if (pBatch == NULL) return 0;
  return rocksdb_writebatch_count(pBatch);
}

void    streamStateClearBatch(void* pBatch) { rocksdb_writebatch_clear((rocksdb_writebatch_t*)pBatch); }
void    streamStateDestroyBatch(void* pBatch) { rocksdb_writebatch_destroy((rocksdb_writebatch_t*)pBatch); }
int32_t streamStatePutBatch(SStreamState* pState, const char* cfName, rocksdb_writebatch_t* pBatch, void* key,
dengyihao's avatar
dengyihao 已提交
2114
                            void* val, int32_t vlen, int64_t ttl) {
Y
yihaoDeng 已提交
2115 2116
  SBackendCfWrapper* wrapper = pState->pTdbState->pBackendCfWrapper;
  int                i = streamStateGetCfIdx(pState, cfName);
dengyihao's avatar
dengyihao 已提交
2117 2118 2119 2120 2121 2122 2123 2124 2125

  if (i < 0) {
    qError("streamState failed to put to cf name:%s", cfName);
    return -1;
  }
  char    buf[128] = {0};
  int32_t klen = ginitDict[i].enFunc((void*)key, buf);

  char*                           ttlV = NULL;
dengyihao's avatar
dengyihao 已提交
2126
  int32_t                         ttlVLen = ginitDict[i].enValueFunc(val, vlen, ttl, &ttlV);
Y
yihaoDeng 已提交
2127
  rocksdb_column_family_handle_t* pCf = wrapper->pHandle[ginitDict[i].idx];
dengyihao's avatar
dengyihao 已提交
2128 2129 2130 2131
  rocksdb_writebatch_put_cf((rocksdb_writebatch_t*)pBatch, pCf, buf, (size_t)klen, ttlV, (size_t)ttlVLen);
  taosMemoryFree(ttlV);
  return 0;
}
dengyihao's avatar
dengyihao 已提交
2132 2133 2134 2135 2136 2137 2138
int32_t streamStatePutBatchOptimize(SStreamState* pState, int32_t cfIdx, rocksdb_writebatch_t* pBatch, void* key,
                                    void* val, int32_t vlen, int64_t ttl, void* tmpBuf) {
  char    buf[128] = {0};
  int32_t klen = ginitDict[cfIdx].enFunc((void*)key, buf);
  char*   ttlV = tmpBuf;
  int32_t ttlVLen = ginitDict[cfIdx].enValueFunc(val, vlen, ttl, &ttlV);

Y
yihaoDeng 已提交
2139
  SBackendCfWrapper* wrapper = pState->pTdbState->pBackendCfWrapper;
Y
yihaoDeng 已提交
2140 2141

  rocksdb_column_family_handle_t* pCf = wrapper->pHandle[ginitDict[cfIdx].idx];
dengyihao's avatar
dengyihao 已提交
2142 2143 2144 2145 2146 2147 2148
  rocksdb_writebatch_put_cf((rocksdb_writebatch_t*)pBatch, pCf, buf, (size_t)klen, ttlV, (size_t)ttlVLen);

  if (tmpBuf == NULL) {
    taosMemoryFree(ttlV);
  }
  return 0;
}
dengyihao's avatar
dengyihao 已提交
2149
int32_t streamStatePutBatch_rocksdb(SStreamState* pState, void* pBatch) {
Y
yihaoDeng 已提交
2150 2151
  char*              err = NULL;
  SBackendCfWrapper* wrapper = pState->pTdbState->pBackendCfWrapper;
Y
yihaoDeng 已提交
2152
  rocksdb_write(wrapper->rocksdb, wrapper->writeOpts, (rocksdb_writebatch_t*)pBatch, &err);
dengyihao's avatar
dengyihao 已提交
2153 2154 2155 2156 2157 2158
  if (err != NULL) {
    qError("streamState failed to write batch, err:%s", err);
    taosMemoryFree(err);
    return -1;
  }
  return 0;
L
liuyao 已提交
2159
}
Y
yihaoDeng 已提交
2160 2161 2162 2163 2164 2165 2166 2167 2168 2169

uint32_t nextPow2(uint32_t x) {
  x = x - 1;
  x = x | (x >> 1);
  x = x | (x >> 2);
  x = x | (x >> 4);
  x = x | (x >> 8);
  x = x | (x >> 16);
  return x + 1;
}