crash_gen.py 78.3 KB
Newer Older
S
Shuduo Sang 已提交
1
# -----!/usr/bin/python3.7
S
Steven Li 已提交
2 3 4 5 6 7 8 9 10 11 12 13
###################################################################
#           Copyright (c) 2016 by TAOS Technologies, Inc.
#                     All rights reserved.
#
#  This file is proprietary and confidential to TAOS Technologies.
#  No part of this file may be reproduced, stored, transmitted,
#  disclosed or used in any form or by any means other than as
#  expressly provided by the written permission from Jianhui Tao
#
###################################################################

# -*- coding: utf-8 -*-
S
Shuduo Sang 已提交
14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38
# For type hinting before definition, ref:
# https://stackoverflow.com/questions/33533148/how-do-i-specify-that-the-return-type-of-a-method-is-the-same-as-the-class-itsel
from __future__ import annotations
import taos
import crash_gen
from util.sql import *
from util.cases import *
from util.dnodes import *
from util.log import *
from queue import Queue, Empty
from typing import IO
from typing import Set
from typing import Dict
from typing import List
from requests.auth import HTTPBasicAuth
import textwrap
import datetime
import logging
import time
import random
import threading
import requests
import copy
import argparse
import getopt
39

S
Steven Li 已提交
40
import sys
41
import os
42 43
import io
import signal
44
import traceback
45 46 47 48
# Require Python 3
if sys.version_info[0] < 3:
    raise Exception("Must be using Python 3")

S
Steven Li 已提交
49

S
Shuduo Sang 已提交
50
# Global variables, tried to keep a small number.
51 52 53

# Command-line/Environment Configurations, will set a bit later
# ConfigNameSpace = argparse.Namespace
S
Shuduo Sang 已提交
54
gConfig = argparse.Namespace()  # Dummy value, will be replaced later
55
logger = None
S
Steven Li 已提交
56

S
Shuduo Sang 已提交
57 58

def runThread(wt: WorkerThread):
59
    wt.run()
60

S
Shuduo Sang 已提交
61

62 63
class CrashGenError(Exception):
    def __init__(self, msg=None, errno=None):
S
Shuduo Sang 已提交
64
        self.msg = msg
65
        self.errno = errno
S
Shuduo Sang 已提交
66

67 68 69
    def __str__(self):
        return self.msg

S
Shuduo Sang 已提交
70

S
Steven Li 已提交
71
class WorkerThread:
S
Shuduo Sang 已提交
72 73 74 75 76
    def __init__(self, pool: ThreadPool, tid,
                 tc: ThreadCoordinator,
                 # te: TaskExecutor,
                 ):  # note: main thread context!
        # self._curStep = -1
77
        self._pool = pool
S
Shuduo Sang 已提交
78 79
        self._tid = tid
        self._tc = tc  # type: ThreadCoordinator
S
Steven Li 已提交
80
        # self.threadIdent = threading.get_ident()
81 82
        self._thread = threading.Thread(target=runThread, args=(self,))
        self._stepGate = threading.Event()
S
Steven Li 已提交
83

84
        # Let us have a DB connection of our own
S
Shuduo Sang 已提交
85
        if (gConfig.per_thread_db_connection):  # type: ignore
86
            # print("connector_type = {}".format(gConfig.connector_type))
S
Shuduo Sang 已提交
87 88
            self._dbConn = DbConn.createNative() if (
                gConfig.connector_type == 'native') else DbConn.createRest()
89

S
Shuduo Sang 已提交
90
        self._dbInUse = False  # if "use db" was executed already
91

92
    def logDebug(self, msg):
S
Steven Li 已提交
93
        logger.debug("    TRD[{}] {}".format(self._tid, msg))
94 95

    def logInfo(self, msg):
S
Steven Li 已提交
96
        logger.info("    TRD[{}] {}".format(self._tid, msg))
97

98 99 100 101
    def dbInUse(self):
        return self._dbInUse

    def useDb(self):
S
Shuduo Sang 已提交
102
        if (not self._dbInUse):
103 104 105
            self.execSql("use db")
        self._dbInUse = True

106
    def getTaskExecutor(self):
S
Shuduo Sang 已提交
107
        return self._tc.getTaskExecutor()
108

S
Steven Li 已提交
109
    def start(self):
110
        self._thread.start()  # AFTER the thread is recorded
S
Steven Li 已提交
111

S
Shuduo Sang 已提交
112
    def run(self):
S
Steven Li 已提交
113
        # initialization after thread starts, in the thread context
114
        # self.isSleeping = False
115 116
        logger.info("Starting to run thread: {}".format(self._tid))

S
Shuduo Sang 已提交
117
        if (gConfig.per_thread_db_connection):  # type: ignore
118
            logger.debug("Worker thread openning database connection")
119
            self._dbConn.open()
S
Steven Li 已提交
120

S
Shuduo Sang 已提交
121 122
        self._doTaskLoop()

123
        # clean up
S
Shuduo Sang 已提交
124
        if (gConfig.per_thread_db_connection):  # type: ignore
125
            self._dbConn.close()
126

S
Shuduo Sang 已提交
127
    def _doTaskLoop(self):
128 129
        # while self._curStep < self._pool.maxSteps:
        # tc = ThreadCoordinator(None)
S
Shuduo Sang 已提交
130 131
        while True:
            tc = self._tc  # Thread Coordinator, the overall master
132
            tc.crossStepBarrier()  # shared barrier first, INCLUDING the last one
S
Shuduo Sang 已提交
133 134 135
            logger.debug(
                "[TRD] Worker thread [{}] exited barrier...".format(
                    self._tid))
136
            self.crossStepGate()   # then per-thread gate, after being tapped
S
Shuduo Sang 已提交
137 138 139
            logger.debug(
                "[TRD] Worker thread [{}] exited step gate...".format(
                    self._tid))
140
            if not self._tc.isRunning():
S
Shuduo Sang 已提交
141 142
                logger.debug(
                    "[TRD] Thread Coordinator not running any more, worker thread now stopping...")
143 144
                break

145
            # Fetch a task from the Thread Coordinator
S
Shuduo Sang 已提交
146 147 148
            logger.debug(
                "[TRD] Worker thread [{}] about to fetch task".format(
                    self._tid))
149
            task = tc.fetchTask()
150 151

            # Execute such a task
S
Shuduo Sang 已提交
152 153 154
            logger.debug(
                "[TRD] Worker thread [{}] about to execute task: {}".format(
                    self._tid, task.__class__.__name__))
155
            task.execute(self)
156
            tc.saveExecutedTask(task)
S
Shuduo Sang 已提交
157 158 159 160 161
            logger.debug(
                "[TRD] Worker thread [{}] finished executing task".format(
                    self._tid))

            self._dbInUse = False  # there may be changes between steps
162

S
Shuduo Sang 已提交
163 164
    def verifyThreadSelf(self):  # ensure we are called by this own thread
        if (threading.get_ident() != self._thread.ident):
S
Steven Li 已提交
165 166
            raise RuntimeError("Unexpectly called from other threads")

S
Shuduo Sang 已提交
167 168
    def verifyThreadMain(self):  # ensure we are called by the main thread
        if (threading.get_ident() != threading.main_thread().ident):
S
Steven Li 已提交
169 170 171
            raise RuntimeError("Unexpectly called from other threads")

    def verifyThreadAlive(self):
S
Shuduo Sang 已提交
172
        if (not self._thread.is_alive()):
S
Steven Li 已提交
173 174
            raise RuntimeError("Unexpected dead thread")

175
    # A gate is different from a barrier in that a thread needs to be "tapped"
S
Steven Li 已提交
176 177
    def crossStepGate(self):
        self.verifyThreadAlive()
S
Shuduo Sang 已提交
178 179
        self.verifyThreadSelf()  # only allowed by ourselves

180
        # Wait again at the "gate", waiting to be "tapped"
S
Shuduo Sang 已提交
181 182 183 184
        logger.debug(
            "[TRD] Worker thread {} about to cross the step gate".format(
                self._tid))
        self._stepGate.wait()
185
        self._stepGate.clear()
S
Shuduo Sang 已提交
186

187
        # self._curStep += 1  # off to a new step...
S
Steven Li 已提交
188

S
Shuduo Sang 已提交
189
    def tapStepGate(self):  # give it a tap, release the thread waiting there
S
Steven Li 已提交
190
        self.verifyThreadAlive()
S
Shuduo Sang 已提交
191 192
        self.verifyThreadMain()  # only allowed for main thread

S
Steven Li 已提交
193
        logger.debug("[TRD] Tapping worker thread {}".format(self._tid))
S
Shuduo Sang 已提交
194 195
        self._stepGate.set()  # wake up!
        time.sleep(0)  # let the released thread run a bit
196

S
Shuduo Sang 已提交
197 198 199
    def execSql(self, sql):  # TODO: expose DbConn directly
        if (gConfig.per_thread_db_connection):
            return self._dbConn.execute(sql)
200
        else:
201
            return self._tc.getDbManager().getDbConn().execute(sql)
202

S
Shuduo Sang 已提交
203 204 205
    def querySql(self, sql):  # TODO: expose DbConn directly
        if (gConfig.per_thread_db_connection):
            return self._dbConn.query(sql)
206 207 208 209
        else:
            return self._tc.getDbManager().getDbConn().query(sql)

    def getQueryResult(self):
S
Shuduo Sang 已提交
210 211
        if (gConfig.per_thread_db_connection):
            return self._dbConn.getQueryResult()
212 213 214
        else:
            return self._tc.getDbManager().getDbConn().getQueryResult()

215
    def getDbConn(self):
S
Shuduo Sang 已提交
216 217
        if (gConfig.per_thread_db_connection):
            return self._dbConn
218
        else:
219
            return self._tc.getDbManager().getDbConn()
220

221 222
    # def querySql(self, sql): # not "execute", since we are out side the DB context
    #     if ( gConfig.per_thread_db_connection ):
S
Shuduo Sang 已提交
223
    #         return self._dbConn.query(sql)
224 225
    #     else:
    #         return self._tc.getDbState().getDbConn().query(sql)
226

227
# The coordinator of all worker threads, mostly running in main thread
S
Shuduo Sang 已提交
228 229


230
class ThreadCoordinator:
231
    def __init__(self, pool: ThreadPool, dbManager):
S
Shuduo Sang 已提交
232
        self._curStep = -1  # first step is 0
233
        self._pool = pool
234
        # self._wd = wd
S
Shuduo Sang 已提交
235
        self._te = None  # prepare for every new step
236
        self._dbManager = dbManager
S
Shuduo Sang 已提交
237 238
        self._executedTasks: List[Task] = []  # in a given step
        self._lock = threading.RLock()  # sync access for a few things
S
Steven Li 已提交
239

S
Shuduo Sang 已提交
240 241
        self._stepBarrier = threading.Barrier(
            self._pool.numThreads + 1)  # one barrier for all threads
242
        self._execStats = ExecutionStats()
243
        self._runStatus = MainExec.STATUS_RUNNING
S
Steven Li 已提交
244

245 246 247
    def getTaskExecutor(self):
        return self._te

S
Shuduo Sang 已提交
248
    def getDbManager(self) -> DbManager:
249
        return self._dbManager
250

251 252 253
    def crossStepBarrier(self):
        self._stepBarrier.wait()

254 255 256 257
    def requestToStop(self):
        self._runStatus = MainExec.STATUS_STOPPING
        self._execStats.registerFailure("User Interruption")

S
Shuduo Sang 已提交
258
    def run(self):
259
        self._pool.createAndStartThreads(self)
S
Steven Li 已提交
260 261

        # Coordinate all threads step by step
S
Shuduo Sang 已提交
262 263 264
        self._curStep = -1  # not started yet
        maxSteps = gConfig.max_steps  # type: ignore
        self._execStats.startExec()  # start the stop watch
265 266
        transitionFailed = False
        hasAbortedTask = False
S
Shuduo Sang 已提交
267 268 269 270 271 272 273 274
        while(self._curStep < maxSteps - 1 and
              (not transitionFailed) and
              (self._runStatus == MainExec.STATUS_RUNNING) and
                (not hasAbortedTask)):  # maxStep==10, last curStep should be 9

            if not gConfig.debug:
                # print this only if we are not in debug mode
                print(".", end="", flush=True)
S
Steven Li 已提交
275
            logger.debug("[TRD] Main thread going to sleep")
276

277
            # Now main thread (that's us) is ready to enter a step
S
Shuduo Sang 已提交
278 279 280 281
            # let other threads go past the pool barrier, but wait at the
            # thread gate
            self.crossStepBarrier()
            self._stepBarrier.reset()  # Other worker threads should now be at the "gate"
282 283

            # At this point, all threads should be pass the overall "barrier" and before the per-thread "gate"
S
Shuduo Sang 已提交
284 285
            # We use this period to do house keeping work, when all worker
            # threads are QUIET.
286
            hasAbortedTask = False
S
Shuduo Sang 已提交
287 288
            for task in self._executedTasks:
                if task.isAborted():
289 290 291 292
                    print("Task aborted: {}".format(task))
                    hasAbortedTask = True
                    break

S
Shuduo Sang 已提交
293
            if hasAbortedTask:  # do transition only if tasks are error free
294
                self._execStats.registerFailure("Aborted Task Encountered")
S
Shuduo Sang 已提交
295
            else:
296 297 298
                try:
                    sm = self._dbManager.getStateMachine()
                    logger.debug("[STT] starting transitions")
S
Shuduo Sang 已提交
299 300
                    # at end of step, transiton the DB state
                    sm.transition(self._executedTasks)
301
                    logger.debug("[STT] transition ended")
S
Shuduo Sang 已提交
302 303 304
                    # Due to limitation (or maybe not) of the Python library,
                    # we cannot share connections across threads
                    if sm.hasDatabase():
305 306 307
                        for t in self._pool.threadList:
                            logger.debug("[DB] use db for all worker threads")
                            t.useDb()
S
Shuduo Sang 已提交
308 309
                            # t.execSql("use db") # main thread executing "use
                            # db" on behalf of every worker thread
310
                except taos.error.ProgrammingError as err:
S
Shuduo Sang 已提交
311
                    if (err.msg == 'network unavailable'):  # broken DB connection
312 313 314
                        logger.info("DB connection broken, execution failed")
                        traceback.print_stack()
                        transitionFailed = True
S
Shuduo Sang 已提交
315
                        self._te = None  # Not running any more
316
                        self._execStats.registerFailure("Broken DB Connection")
S
Shuduo Sang 已提交
317 318
                        # continue # don't do that, need to tap all threads at
                        # end, and maybe signal them to stop
319
                    else:
S
Shuduo Sang 已提交
320
                        raise
321 322
            # finally:
            #     pass
S
Shuduo Sang 已提交
323 324

            self.resetExecutedTasks()  # clear the tasks after we are done
325 326

            # Get ready for next step
S
Steven Li 已提交
327
            logger.debug("<-- Step {} finished".format(self._curStep))
S
Shuduo Sang 已提交
328 329 330 331
            self._curStep += 1  # we are about to get into next step. TODO: race condition here!
            # Now not all threads had time to go to sleep
            logger.debug(
                "\r\n\n--> Step {} starts with main thread waking up".format(self._curStep))
332

333
            # A new TE for the new step
S
Shuduo Sang 已提交
334
            if not transitionFailed:  # only if not failed
335
                self._te = TaskExecutor(self._curStep)
336

S
Shuduo Sang 已提交
337 338 339 340 341 342
            logger.debug(
                "[TRD] Main thread waking up at step {}, tapping worker threads".format(
                    self._curStep))  # Now not all threads had time to go to sleep
            # Worker threads will wake up at this point, and each execute it's
            # own task
            self.tapAllThreads()
S
Steven Li 已提交
343

344
        logger.debug("Main thread ready to finish up...")
S
Shuduo Sang 已提交
345 346
        if not transitionFailed:  # only in regular situations
            self.crossStepBarrier()  # Cross it one last time, after all threads finish
347 348
            self._stepBarrier.reset()
            logger.debug("Main thread in exclusive zone...")
S
Shuduo Sang 已提交
349
            self._te = None  # No more executor, time to end
350
            logger.debug("Main thread tapping all threads one last time...")
S
Shuduo Sang 已提交
351
            self.tapAllThreads()  # Let the threads run one last time
352

353
        logger.debug("Main thread joining all threads")
S
Shuduo Sang 已提交
354
        self._pool.joinAll()  # Get all threads to finish
355
        logger.info("\nAll worker threads finished")
356 357
        self._execStats.endExec()

358 359
    def printStats(self):
        self._execStats.printStats()
S
Steven Li 已提交
360

S
Steven Li 已提交
361 362 363 364 365 366
    def isFailed(self):
        return self._execStats.isFailed()

    def getExecStats(self):
        return self._execStats

S
Shuduo Sang 已提交
367
    def tapAllThreads(self):  # in a deterministic manner
S
Steven Li 已提交
368
        wakeSeq = []
S
Shuduo Sang 已提交
369 370
        for i in range(self._pool.numThreads):  # generate a random sequence
            if Dice.throw(2) == 1:
S
Steven Li 已提交
371 372 373
                wakeSeq.append(i)
            else:
                wakeSeq.insert(0, i)
S
Shuduo Sang 已提交
374 375 376
        logger.debug(
            "[TRD] Main thread waking up worker threads: {}".format(
                str(wakeSeq)))
377
        # TODO: set dice seed to a deterministic value
S
Steven Li 已提交
378
        for i in wakeSeq:
S
Shuduo Sang 已提交
379 380 381
            # TODO: maybe a bit too deep?!
            self._pool.threadList[i].tapStepGate()
            time.sleep(0)  # yield
S
Steven Li 已提交
382

383
    def isRunning(self):
S
Shuduo Sang 已提交
384
        return self._te is not None
385

S
Shuduo Sang 已提交
386 387
    def fetchTask(self) -> Task:
        if (not self.isRunning()):  # no task
388
            raise RuntimeError("Cannot fetch task when not running")
389 390
        # return self._wd.pickTask()
        # Alternatively, let's ask the DbState for the appropriate task
391 392 393 394 395 396 397
        # dbState = self.getDbState()
        # tasks = dbState.getTasksAtState() # TODO: create every time?
        # nTasks = len(tasks)
        # i = Dice.throw(nTasks)
        # logger.debug(" (dice:{}/{}) ".format(i, nTasks))
        # # return copy.copy(tasks[i]) # Needs a fresh copy, to save execution results, etc.
        # return tasks[i].clone() # TODO: still necessary?
S
Shuduo Sang 已提交
398 399 400 401 402
        # pick a task type for current state
        taskType = self.getDbManager().getStateMachine().pickTaskType()
        return taskType(
            self.getDbManager(),
            self._execStats)  # create a task from it
403 404

    def resetExecutedTasks(self):
S
Shuduo Sang 已提交
405
        self._executedTasks = []  # should be under single thread
406 407 408 409

    def saveExecutedTask(self, task):
        with self._lock:
            self._executedTasks.append(task)
410 411

# We define a class to run a number of threads in locking steps.
S
Shuduo Sang 已提交
412 413


414
class ThreadPool:
415
    def __init__(self, numThreads, maxSteps):
416 417 418 419
        self.numThreads = numThreads
        self.maxSteps = maxSteps
        # Internal class variables
        self.curStep = 0
S
Shuduo Sang 已提交
420 421
        self.threadList = []  # type: List[WorkerThread]

422
    # starting to run all the threads, in locking steps
423
    def createAndStartThreads(self, tc: ThreadCoordinator):
S
Shuduo Sang 已提交
424 425
        for tid in range(0, self.numThreads):  # Create the threads
            workerThread = WorkerThread(self, tid, tc)
426
            self.threadList.append(workerThread)
S
Shuduo Sang 已提交
427
            workerThread.start()  # start, but should block immediately before step 0
428 429 430 431 432 433

    def joinAll(self):
        for workerThread in self.threadList:
            logger.debug("Joining thread...")
            workerThread._thread.join()

434 435
# A queue of continguous POSITIVE integers, used by DbManager to generate continuous numbers
# for new table names
S
Shuduo Sang 已提交
436 437


S
Steven Li 已提交
438 439
class LinearQueue():
    def __init__(self):
440
        self.firstIndex = 1  # 1st ever element
S
Steven Li 已提交
441
        self.lastIndex = 0
S
Shuduo Sang 已提交
442 443
        self._lock = threading.RLock()  # our functions may call each other
        self.inUse = set()  # the indexes that are in use right now
S
Steven Li 已提交
444

445
    def toText(self):
S
Shuduo Sang 已提交
446 447
        return "[{}..{}], in use: {}".format(
            self.firstIndex, self.lastIndex, self.inUse)
448 449

    # Push (add new element, largest) to the tail, and mark it in use
S
Shuduo Sang 已提交
450
    def push(self):
451
        with self._lock:
S
Shuduo Sang 已提交
452 453
            # if ( self.isEmpty() ):
            #     self.lastIndex = self.firstIndex
454
            #     return self.firstIndex
455 456
            # Otherwise we have something
            self.lastIndex += 1
457 458
            self.allocate(self.lastIndex)
            # self.inUse.add(self.lastIndex) # mark it in use immediately
459
            return self.lastIndex
S
Steven Li 已提交
460 461

    def pop(self):
462
        with self._lock:
S
Shuduo Sang 已提交
463 464 465 466
            if (self.isEmpty()):
                # raise RuntimeError("Cannot pop an empty queue")
                return False  # TODO: None?

467
            index = self.firstIndex
S
Shuduo Sang 已提交
468
            if (index in self.inUse):
469 470
                return False

471 472 473 474 475 476 477
            self.firstIndex += 1
            return index

    def isEmpty(self):
        return self.firstIndex > self.lastIndex

    def popIfNotEmpty(self):
478
        with self._lock:
479 480 481 482
            if (self.isEmpty()):
                return 0
            return self.pop()

S
Steven Li 已提交
483
    def allocate(self, i):
484
        with self._lock:
485
            # logger.debug("LQ allocating item {}".format(i))
S
Shuduo Sang 已提交
486 487 488
            if (i in self.inUse):
                raise RuntimeError(
                    "Cannot re-use same index in queue: {}".format(i))
489 490
            self.inUse.add(i)

S
Steven Li 已提交
491
    def release(self, i):
492
        with self._lock:
493
            # logger.debug("LQ releasing item {}".format(i))
S
Shuduo Sang 已提交
494
            self.inUse.remove(i)  # KeyError possible, TODO: why?
495 496 497 498

    def size(self):
        return self.lastIndex + 1 - self.firstIndex

S
Steven Li 已提交
499
    def pickAndAllocate(self):
S
Shuduo Sang 已提交
500
        if (self.isEmpty()):
501 502
            return None
        with self._lock:
S
Shuduo Sang 已提交
503
            cnt = 0  # counting the interations
504 505
            while True:
                cnt += 1
S
Shuduo Sang 已提交
506
                if (cnt > self.size() * 10):  # 10x iteration already
507 508
                    # raise RuntimeError("Failed to allocate LinearQueue element")
                    return None
S
Shuduo Sang 已提交
509 510
                ret = Dice.throwRange(self.firstIndex, self.lastIndex + 1)
                if (ret not in self.inUse):
511 512 513
                    self.allocate(ret)
                    return ret

S
Shuduo Sang 已提交
514

515
class DbConn:
516 517 518 519 520 521 522 523 524 525 526
    TYPE_NATIVE = "native-c"
    TYPE_REST = "rest-api"
    TYPE_INVALID = "invalid"

    @classmethod
    def create(cls, connType):
        if connType == cls.TYPE_NATIVE:
            return DbConnNative()
        elif connType == cls.TYPE_REST:
            return DbConnRest()
        else:
S
Shuduo Sang 已提交
527 528
            raise RuntimeError(
                "Unexpected connection type: {}".format(connType))
529 530 531 532 533 534 535 536 537

    @classmethod
    def createNative(cls):
        return cls.create(cls.TYPE_NATIVE)

    @classmethod
    def createRest(cls):
        return cls.create(cls.TYPE_REST)

538 539
    def __init__(self):
        self.isOpen = False
540 541 542
        self._type = self.TYPE_INVALID

    def open(self):
S
Shuduo Sang 已提交
543
        if (self.isOpen):
544 545
            raise RuntimeError("Cannot re-open an existing DB connection")

546 547
        # below implemented by child classes
        self.openByType()
548

S
Shuduo Sang 已提交
549 550 551
        logger.debug(
            "[DB] data connection opened, type = {}".format(
                self._type))
552 553
        self.isOpen = True

S
Shuduo Sang 已提交
554 555 556 557
    def resetDb(self):  # reset the whole database, etc.
        if (not self.isOpen):
            raise RuntimeError(
                "Cannot reset database until connection is open")
558 559
        # self._tdSql.prepare() # Recreate database, etc.

560
        self.execute('drop database if exists db')
561 562
        logger.debug("Resetting DB, dropped database")
        # self._cursor.execute('create database db')
563
        # self._cursor.execute('use db')
564 565
        # tdSql.execute('show databases')

S
Shuduo Sang 已提交
566
    def queryScalar(self, sql) -> int:
567 568
        return self._queryAny(sql)

S
Shuduo Sang 已提交
569
    def queryString(self, sql) -> str:
570 571
        return self._queryAny(sql)

S
Shuduo Sang 已提交
572 573 574 575
    def _queryAny(self, sql):  # actual query result as an int
        if (not self.isOpen):
            raise RuntimeError(
                "Cannot query database until connection is open")
576
        nRows = self.query(sql)
S
Shuduo Sang 已提交
577 578 579 580
        if nRows != 1:
            raise RuntimeError(
                "Unexpected result for query: {}, rows = {}".format(
                    sql, nRows))
581
        if self.getResultRows() != 1 or self.getResultCols() != 1:
S
Shuduo Sang 已提交
582 583
            raise RuntimeError(
                "Unexpected result set for query: {}".format(sql))
584 585 586 587
        return self.getQueryResult()[0][0]

    def execute(self, sql):
        raise RuntimeError("Unexpected execution, should be overriden")
S
Shuduo Sang 已提交
588

589 590
    def openByType(self):
        raise RuntimeError("Unexpected execution, should be overriden")
S
Shuduo Sang 已提交
591

592 593
    def getQueryResult(self):
        raise RuntimeError("Unexpected execution, should be overriden")
S
Shuduo Sang 已提交
594

595 596
    def getResultRows(self):
        raise RuntimeError("Unexpected execution, should be overriden")
S
Shuduo Sang 已提交
597

598 599 600 601
    def getResultCols(self):
        raise RuntimeError("Unexpected execution, should be overriden")

# Sample: curl -u root:taosdata -d "show databases" localhost:6020/rest/sql
S
Shuduo Sang 已提交
602 603


604 605 606 607
class DbConnRest(DbConn):
    def __init__(self):
        super().__init__()
        self._type = self.TYPE_REST
S
Shuduo Sang 已提交
608
        self._url = "http://localhost:6020/rest/sql"  # fixed for now
609 610
        self._result = None

S
Shuduo Sang 已提交
611 612 613
    def openByType(self):  # Open connection
        pass  # do nothing, always open

614
    def close(self):
S
Shuduo Sang 已提交
615 616 617
        if (not self.isOpen):
            raise RuntimeError(
                "Cannot clean up database until connection is open")
618 619 620 621 622
        # Do nothing for REST
        logger.debug("[DB] REST Database connection closed")
        self.isOpen = False

    def _doSql(self, sql):
S
Shuduo Sang 已提交
623 624 625
        r = requests.post(self._url,
                          data=sql,
                          auth=HTTPBasicAuth('root', 'taosdata'))
626 627
        rj = r.json()
        # Sanity check for the "Json Result"
S
Shuduo Sang 已提交
628
        if ('status' not in rj):
629 630
            raise RuntimeError("No status in REST response")

S
Shuduo Sang 已提交
631 632 633 634
        if rj['status'] == 'error':  # clearly reported error
            if ('code' not in rj):  # error without code
                raise RuntimeError("REST error return without code")
            errno = rj['code']  # May need to massage this in the future
635
            # print("Raising programming error with REST return: {}".format(rj))
S
Shuduo Sang 已提交
636 637
            raise taos.error.ProgrammingError(
                rj['desc'], errno)  # todo: check existance of 'desc'
638

S
Shuduo Sang 已提交
639 640 641 642
        if rj['status'] != 'succ':  # better be this
            raise RuntimeError(
                "Unexpected REST return status: {}".format(
                    rj['status']))
643 644

        nRows = rj['rows'] if ('rows' in rj) else 0
S
Shuduo Sang 已提交
645
        self._result = rj
646 647
        return nRows

S
Shuduo Sang 已提交
648 649 650 651
    def execute(self, sql):
        if (not self.isOpen):
            raise RuntimeError(
                "Cannot execute database commands until connection is open")
652 653
        logger.debug("[SQL-REST] Executing SQL: {}".format(sql))
        nRows = self._doSql(sql)
S
Shuduo Sang 已提交
654 655
        logger.debug(
            "[SQL-REST] Execution Result, nRows = {}, SQL = {}".format(nRows, sql))
656 657
        return nRows

S
Shuduo Sang 已提交
658
    def query(self, sql):  # return rows affected
659 660 661 662 663 664 665 666 667 668 669 670 671
        return self.execute(sql)

    def getQueryResult(self):
        return self._result['data']

    def getResultRows(self):
        print(self._result)
        raise RuntimeError("TBD")
        # return self._tdSql.queryRows

    def getResultCols(self):
        print(self._result)
        raise RuntimeError("TBD")
S
Shuduo Sang 已提交
672 673


674 675 676 677
class DbConnNative(DbConn):
    def __init__(self):
        super().__init__()
        self._type = self.TYPE_REST
S
Shuduo Sang 已提交
678
        self._conn = None
679
        self._cursor = None
S
Shuduo Sang 已提交
680

681 682 683 684 685 686 687 688 689 690 691 692 693 694 695
    def getBuildPath(self):
        selfPath = os.path.dirname(os.path.realpath(__file__))
        if ("community" in selfPath):
            projPath = selfPath[:selfPath.find("communit")]
        else:
            projPath = selfPath[:selfPath.find("tests")]

        for root, dirs, files in os.walk(projPath):
            if ("taosd" in files):
                rootRealPath = os.path.dirname(os.path.realpath(root))
                if ("packaging" not in rootRealPath):
                    buildPath = root[:root.find("build")]
                    break
        return buildPath

S
Shuduo Sang 已提交
696
    def openByType(self):  # Open connection
697 698 699
#        cfgPath = "../../build/test/cfg"
        cfgPath = self.getBuildPath() + "/test/cfg"
        print("CBD: cfgPath=%s" % cfgPath)
S
Shuduo Sang 已提交
700 701 702
        self._conn = taos.connect(
            host="127.0.0.1",
            config=cfgPath)  # TODO: make configurable
703 704 705 706
        self._cursor = self._conn.cursor()

        # Get the connection/cursor ready
        self._cursor.execute('reset query cache')
S
Shuduo Sang 已提交
707 708
        # self._cursor.execute('use db') # do this at the beginning of every
        # step
709 710 711 712

        # Open connection
        self._tdSql = TDSql()
        self._tdSql.init(self._cursor)
S
Shuduo Sang 已提交
713

714
    def close(self):
S
Shuduo Sang 已提交
715 716 717
        if (not self.isOpen):
            raise RuntimeError(
                "Cannot clean up database until connection is open")
718
        self._tdSql.close()
719
        logger.debug("[DB] Database connection closed")
720
        self.isOpen = False
S
Steven Li 已提交
721

S
Shuduo Sang 已提交
722 723 724 725
    def execute(self, sql):
        if (not self.isOpen):
            raise RuntimeError(
                "Cannot execute database commands until connection is open")
726 727
        logger.debug("[SQL] Executing SQL: {}".format(sql))
        nRows = self._tdSql.execute(sql)
S
Shuduo Sang 已提交
728 729 730
        logger.debug(
            "[SQL] Execution Result, nRows = {}, SQL = {}".format(
                nRows, sql))
731
        return nRows
S
Steven Li 已提交
732

S
Shuduo Sang 已提交
733 734 735 736
    def query(self, sql):  # return rows affected
        if (not self.isOpen):
            raise RuntimeError(
                "Cannot query database until connection is open")
737 738
        logger.debug("[SQL] Executing SQL: {}".format(sql))
        nRows = self._tdSql.query(sql)
S
Shuduo Sang 已提交
739 740 741
        logger.debug(
            "[SQL] Query Result, nRows = {}, SQL = {}".format(
                nRows, sql))
742
        return nRows
743
        # results are in: return self._tdSql.queryResult
744

745 746 747
    def getQueryResult(self):
        return self._tdSql.queryResult

748 749
    def getResultRows(self):
        return self._tdSql.queryRows
750

751 752
    def getResultCols(self):
        return self._tdSql.queryCols
753

S
Shuduo Sang 已提交
754

755
class AnyState:
S
Shuduo Sang 已提交
756 757 758
    STATE_INVALID = -1
    STATE_EMPTY = 0  # nothing there, no even a DB
    STATE_DB_ONLY = 1  # we have a DB, but nothing else
759
    STATE_TABLE_ONLY = 2  # we have a table, but totally empty
S
Shuduo Sang 已提交
760
    STATE_HAS_DATA = 3  # we have some data in the table
761 762 763 764 765
    _stateNames = ["Invalid", "Empty", "DB_Only", "Table_Only", "Has_Data"]

    STATE_VAL_IDX = 0
    CAN_CREATE_DB = 1
    CAN_DROP_DB = 2
766 767
    CAN_CREATE_FIXED_SUPER_TABLE = 3
    CAN_DROP_FIXED_SUPER_TABLE = 4
768 769 770 771 772 773 774
    CAN_ADD_DATA = 5
    CAN_READ_DATA = 6

    def __init__(self):
        self._info = self.getInfo()

    def __str__(self):
S
Shuduo Sang 已提交
775 776
        # -1 hack to accomodate the STATE_INVALID case
        return self._stateNames[self._info[self.STATE_VAL_IDX] + 1]
777 778 779 780

    def getInfo(self):
        raise RuntimeError("Must be overriden by child classes")

S
Steven Li 已提交
781 782 783 784 785 786
    def equals(self, other):
        if isinstance(other, int):
            return self.getValIndex() == other
        elif isinstance(other, AnyState):
            return self.getValIndex() == other.getValIndex()
        else:
S
Shuduo Sang 已提交
787 788 789
            raise RuntimeError(
                "Unexpected comparison, type = {}".format(
                    type(other)))
S
Steven Li 已提交
790

791 792 793
    def verifyTasksToState(self, tasks, newState):
        raise RuntimeError("Must be overriden by child classes")

S
Steven Li 已提交
794 795 796
    def getValIndex(self):
        return self._info[self.STATE_VAL_IDX]

797 798
    def getValue(self):
        return self._info[self.STATE_VAL_IDX]
S
Shuduo Sang 已提交
799

800 801
    def canCreateDb(self):
        return self._info[self.CAN_CREATE_DB]
S
Shuduo Sang 已提交
802

803 804
    def canDropDb(self):
        return self._info[self.CAN_DROP_DB]
S
Shuduo Sang 已提交
805

806 807
    def canCreateFixedSuperTable(self):
        return self._info[self.CAN_CREATE_FIXED_SUPER_TABLE]
S
Shuduo Sang 已提交
808

809 810
    def canDropFixedSuperTable(self):
        return self._info[self.CAN_DROP_FIXED_SUPER_TABLE]
S
Shuduo Sang 已提交
811

812 813
    def canAddData(self):
        return self._info[self.CAN_ADD_DATA]
S
Shuduo Sang 已提交
814

815 816 817 818 819
    def canReadData(self):
        return self._info[self.CAN_READ_DATA]

    def assertAtMostOneSuccess(self, tasks, cls):
        sCnt = 0
S
Shuduo Sang 已提交
820
        for task in tasks:
821 822 823
            if not isinstance(task, cls):
                continue
            if task.isSuccess():
S
Steven Li 已提交
824
                # task.logDebug("Task success found")
825
                sCnt += 1
S
Shuduo Sang 已提交
826 827 828
                if (sCnt >= 2):
                    raise RuntimeError(
                        "Unexpected more than 1 success with task: {}".format(cls))
829 830 831 832

    def assertIfExistThenSuccess(self, tasks, cls):
        sCnt = 0
        exists = False
S
Shuduo Sang 已提交
833
        for task in tasks:
834 835
            if not isinstance(task, cls):
                continue
S
Shuduo Sang 已提交
836
            exists = True  # we have a valid instance
837 838
            if task.isSuccess():
                sCnt += 1
S
Shuduo Sang 已提交
839 840 841
        if (exists and sCnt <= 0):
            raise RuntimeError(
                "Unexpected zero success for task: {}".format(cls))
842 843

    def assertNoTask(self, tasks, cls):
S
Shuduo Sang 已提交
844
        for task in tasks:
845
            if isinstance(task, cls):
S
Shuduo Sang 已提交
846 847
                raise CrashGenError(
                    "This task: {}, is not expected to be present, given the success/failure of others".format(cls.__name__))
848 849

    def assertNoSuccess(self, tasks, cls):
S
Shuduo Sang 已提交
850
        for task in tasks:
851 852
            if isinstance(task, cls):
                if task.isSuccess():
S
Shuduo Sang 已提交
853 854
                    raise RuntimeError(
                        "Unexpected successful task: {}".format(cls))
855 856

    def hasSuccess(self, tasks, cls):
S
Shuduo Sang 已提交
857
        for task in tasks:
858 859 860 861 862 863
            if not isinstance(task, cls):
                continue
            if task.isSuccess():
                return True
        return False

S
Steven Li 已提交
864
    def hasTask(self, tasks, cls):
S
Shuduo Sang 已提交
865
        for task in tasks:
S
Steven Li 已提交
866 867 868 869
            if isinstance(task, cls):
                return True
        return False

S
Shuduo Sang 已提交
870

871 872 873 874
class StateInvalid(AnyState):
    def getInfo(self):
        return [
            self.STATE_INVALID,
S
Shuduo Sang 已提交
875 876 877
            False, False,  # can create/drop Db
            False, False,  # can create/drop fixed table
            False, False,  # can insert/read data with fixed table
878 879 880 881
        ]

    # def verifyTasksToState(self, tasks, newState):

S
Shuduo Sang 已提交
882

883 884 885 886
class StateEmpty(AnyState):
    def getInfo(self):
        return [
            self.STATE_EMPTY,
S
Shuduo Sang 已提交
887 888 889
            True, False,  # can create/drop Db
            False, False,  # can create/drop fixed table
            False, False,  # can insert/read data with fixed table
890 891
        ]

S
Shuduo Sang 已提交
892 893 894 895 896 897 898
    def verifyTasksToState(self, tasks, newState):
        if (self.hasSuccess(tasks, TaskCreateDb)
                ):  # at EMPTY, if there's succes in creating DB
            if (not self.hasTask(tasks, TaskDropDb)):  # and no drop_db tasks
                # we must have at most one. TODO: compare numbers
                self.assertAtMostOneSuccess(tasks, TaskCreateDb)

899 900 901 902 903 904 905 906 907 908 909

class StateDbOnly(AnyState):
    def getInfo(self):
        return [
            self.STATE_DB_ONLY,
            False, True,
            True, False,
            False, False,
        ]

    def verifyTasksToState(self, tasks, newState):
S
Shuduo Sang 已提交
910 911 912
        if (not self.hasTask(tasks, TaskCreateDb)):
            # only if we don't create any more
            self.assertAtMostOneSuccess(tasks, TaskDropDb)
913
        self.assertIfExistThenSuccess(tasks, TaskDropDb)
S
Steven Li 已提交
914
        # self.assertAtMostOneSuccess(tasks, CreateFixedTableTask) # not true in massively parrallel cases
915
        # Nothing to be said about adding data task
916
        # if ( self.hasSuccess(tasks, DropDbTask) ): # dropped the DB
S
Shuduo Sang 已提交
917 918 919
        # self.assertHasTask(tasks, DropDbTask) # implied by hasSuccess
        # self.assertAtMostOneSuccess(tasks, DropDbTask)
        # self._state = self.STATE_EMPTY
920 921
        # if ( self.hasSuccess(tasks, TaskCreateSuperTable) ): # did not drop db, create table success
        #     # self.assertHasTask(tasks, CreateFixedTableTask) # tried to create table
S
Shuduo Sang 已提交
922
        #     if ( not self.hasTask(tasks, TaskDropSuperTable) ):
923
        #         self.assertAtMostOneSuccess(tasks, TaskCreateSuperTable) # at most 1 attempt is successful, if we don't drop anything
S
Shuduo Sang 已提交
924 925 926 927 928 929
        # self.assertNoTask(tasks, DropDbTask) # should have have tried
        # if ( not self.hasSuccess(tasks, AddFixedDataTask) ): # just created table, no data yet
        #     # can't say there's add-data attempts, since they may all fail
        #     self._state = self.STATE_TABLE_ONLY
        # else:
        #     self._state = self.STATE_HAS_DATA
930 931 932 933
        # What about AddFixedData?
        # elif ( self.hasSuccess(tasks, AddFixedDataTask) ):
        #     self._state = self.STATE_HAS_DATA
        # else: # no success in dropping db tasks, no success in create fixed table? read data should also fail
S
Shuduo Sang 已提交
934
        #     # raise RuntimeError("Unexpected no-success scenario")   # We might just landed all failure tasks,
935 936
        #     self._state = self.STATE_DB_ONLY  # no change

S
Shuduo Sang 已提交
937

938
class StateSuperTableOnly(AnyState):
939 940 941 942 943 944 945 946 947
    def getInfo(self):
        return [
            self.STATE_TABLE_ONLY,
            False, True,
            False, True,
            True, True,
        ]

    def verifyTasksToState(self, tasks, newState):
S
Shuduo Sang 已提交
948 949
        if (self.hasSuccess(tasks, TaskDropSuperTable)
                ):  # we are able to drop the table
950
            #self.assertAtMostOneSuccess(tasks, TaskDropSuperTable)
S
Shuduo Sang 已提交
951 952
            # we must have had recreted it
            self.hasSuccess(tasks, TaskCreateSuperTable)
953

954
            # self._state = self.STATE_DB_ONLY
S
Steven Li 已提交
955 956
        # elif ( self.hasSuccess(tasks, AddFixedDataTask) ): # no success dropping the table, but added data
        #     self.assertNoTask(tasks, DropFixedTableTask) # not true in massively parrallel cases
957
            # self._state = self.STATE_HAS_DATA
S
Steven Li 已提交
958 959 960
        # elif ( self.hasSuccess(tasks, ReadFixedDataTask) ): # no success in prev cases, but was able to read data
            # self.assertNoTask(tasks, DropFixedTableTask)
            # self.assertNoTask(tasks, AddFixedDataTask)
961
            # self._state = self.STATE_TABLE_ONLY # no change
S
Steven Li 已提交
962 963 964
        # else: # did not drop table, did not insert data, did not read successfully, that is impossible
        #     raise RuntimeError("Unexpected no-success scenarios")
        # TODO: need to revamp!!
965

S
Shuduo Sang 已提交
966

967 968 969 970 971 972 973 974 975 976
class StateHasData(AnyState):
    def getInfo(self):
        return [
            self.STATE_HAS_DATA,
            False, True,
            False, True,
            True, True,
        ]

    def verifyTasksToState(self, tasks, newState):
S
Shuduo Sang 已提交
977
        if (newState.equals(AnyState.STATE_EMPTY)):
978
            self.hasSuccess(tasks, TaskDropDb)
S
Shuduo Sang 已提交
979 980 981 982 983 984 985
            if (not self.hasTask(tasks, TaskCreateDb)):
                self.assertAtMostOneSuccess(tasks, TaskDropDb)  # TODO: dicy
        elif (newState.equals(AnyState.STATE_DB_ONLY)):  # in DB only
            if (not self.hasTask(tasks, TaskCreateDb)
                ):  # without a create_db task
                # we must have drop_db task
                self.assertNoTask(tasks, TaskDropDb)
986
            self.hasSuccess(tasks, TaskDropSuperTable)
987
            # self.assertAtMostOneSuccess(tasks, DropFixedSuperTableTask) # TODO: dicy
988 989 990 991
        # elif ( newState.equals(AnyState.STATE_TABLE_ONLY) ): # data deleted
            # self.assertNoTask(tasks, TaskDropDb)
            # self.assertNoTask(tasks, TaskDropSuperTable)
            # self.assertNoTask(tasks, TaskAddData)
S
Steven Li 已提交
992
            # self.hasSuccess(tasks, DeleteDataTasks)
S
Shuduo Sang 已提交
993 994 995 996 997 998 999 1000 1001
        else:  # should be STATE_HAS_DATA
            if (not self.hasTask(tasks, TaskCreateDb)
                ):  # only if we didn't create one
                # we shouldn't have dropped it
                self.assertNoTask(tasks, TaskDropDb)
            if (not self.hasTask(tasks, TaskCreateSuperTable)
                    ):  # if we didn't create the table
                # we should not have a task that drops it
                self.assertNoTask(tasks, TaskDropSuperTable)
1002
            # self.assertIfExistThenSuccess(tasks, ReadFixedDataTask)
S
Steven Li 已提交
1003

S
Shuduo Sang 已提交
1004

1005
class StateMechine:
1006 1007
    def __init__(self, dbConn):
        self._dbConn = dbConn
S
Shuduo Sang 已提交
1008 1009 1010 1011 1012
        self._curState = self._findCurrentState()  # starting state
        # transitition target probabilities, indexed with value of STATE_EMPTY,
        # STATE_DB_ONLY, etc.
        self._stateWeights = [1, 3, 5, 15]

1013 1014 1015
    def getCurrentState(self):
        return self._curState

1016 1017 1018
    def hasDatabase(self):
        return self._curState.canDropDb()  # ha, can drop DB means it has one

1019
    # May be slow, use cautionsly...
S
Shuduo Sang 已提交
1020
    def getTaskTypes(self):  # those that can run (directly/indirectly) from the current state
1021 1022 1023 1024 1025 1026
        def typesToStrings(types):
            ss = []
            for t in types:
                ss.append(t.__name__)
            return ss

S
Shuduo Sang 已提交
1027
        allTaskClasses = StateTransitionTask.__subclasses__()  # all state transition tasks
1028 1029
        firstTaskTypes = []
        for tc in allTaskClasses:
S
Shuduo Sang 已提交
1030
            # t = tc(self) # create task object
1031 1032
            if tc.canBeginFrom(self._curState):
                firstTaskTypes.append(tc)
S
Shuduo Sang 已提交
1033 1034 1035 1036 1037 1038 1039 1040
        # now we have all the tasks that can begin directly from the current
        # state, let's figure out the INDIRECT ones
        taskTypes = firstTaskTypes.copy()  # have to have these
        for task1 in firstTaskTypes:  # each task type gathered so far
            endState = task1.getEndState()  # figure the end state
            if endState is None:  # does not change end state
                continue  # no use, do nothing
            for tc in allTaskClasses:  # what task can further begin from there?
1041
                if tc.canBeginFrom(endState) and (tc not in firstTaskTypes):
S
Shuduo Sang 已提交
1042
                    taskTypes.append(tc)  # gather it
1043 1044

        if len(taskTypes) <= 0:
S
Shuduo Sang 已提交
1045 1046 1047 1048 1049 1050 1051
            raise RuntimeError(
                "No suitable task types found for state: {}".format(
                    self._curState))
        logger.debug(
            "[OPS] Tasks found for state {}: {}".format(
                self._curState,
                typesToStrings(taskTypes)))
1052 1053 1054 1055
        return taskTypes

    def _findCurrentState(self):
        dbc = self._dbConn
S
Shuduo Sang 已提交
1056 1057
        ts = time.time()  # we use this to debug how fast/slow it is to do the various queries to find the current DB state
        if dbc.query("show databases") == 0:  # no database?!
1058
            # logger.debug("Found EMPTY state")
S
Shuduo Sang 已提交
1059 1060 1061
            logger.debug(
                "[STT] empty database found, between {} and {}".format(
                    ts, time.time()))
1062
            return StateEmpty()
S
Shuduo Sang 已提交
1063 1064 1065 1066
        # did not do this when openning connection, and this is NOT the worker
        # thread, which does this on their own
        dbc.execute("use db")
        if dbc.query("show tables") == 0:  # no tables
1067
            # logger.debug("Found DB ONLY state")
S
Shuduo Sang 已提交
1068 1069 1070
            logger.debug(
                "[STT] DB_ONLY found, between {} and {}".format(
                    ts, time.time()))
1071
            return StateDbOnly()
S
Shuduo Sang 已提交
1072 1073
        if dbc.query("SELECT * FROM db.{}".format(DbManager.getFixedSuperTableName())
                     ) == 0:  # no regular tables
1074
            # logger.debug("Found TABLE_ONLY state")
S
Shuduo Sang 已提交
1075 1076 1077
            logger.debug(
                "[STT] SUPER_TABLE_ONLY found, between {} and {}".format(
                    ts, time.time()))
1078
            return StateSuperTableOnly()
S
Shuduo Sang 已提交
1079
        else:  # has actual tables
1080
            # logger.debug("Found HAS_DATA state")
S
Shuduo Sang 已提交
1081 1082 1083
            logger.debug(
                "[STT] HAS_DATA found, between {} and {}".format(
                    ts, time.time()))
1084 1085 1086
            return StateHasData()

    def transition(self, tasks):
S
Shuduo Sang 已提交
1087
        if (len(tasks) == 0):  # before 1st step, or otherwise empty
1088
            logger.debug("[STT] Starting State: {}".format(self._curState))
S
Shuduo Sang 已提交
1089
            return  # do nothing
1090

S
Shuduo Sang 已提交
1091 1092
        # this should show up in the server log, separating steps
        self._dbConn.execute("show dnodes")
1093 1094 1095 1096

        # Generic Checks, first based on the start state
        if self._curState.canCreateDb():
            self._curState.assertIfExistThenSuccess(tasks, TaskCreateDb)
S
Shuduo Sang 已提交
1097 1098
            # self.assertAtMostOneSuccess(tasks, CreateDbTask) # not really, in
            # case of multiple creation and drops
1099 1100 1101

        if self._curState.canDropDb():
            self._curState.assertIfExistThenSuccess(tasks, TaskDropDb)
S
Shuduo Sang 已提交
1102 1103
            # self.assertAtMostOneSuccess(tasks, DropDbTask) # not really in
            # case of drop-create-drop
1104 1105 1106

        # if self._state.canCreateFixedTable():
            # self.assertIfExistThenSuccess(tasks, CreateFixedTableTask) # Not true, DB may be dropped
S
Shuduo Sang 已提交
1107 1108
            # self.assertAtMostOneSuccess(tasks, CreateFixedTableTask) # not
            # really, in case of create-drop-create
1109 1110 1111

        # if self._state.canDropFixedTable():
            # self.assertIfExistThenSuccess(tasks, DropFixedTableTask) # Not True, the whole DB may be dropped
S
Shuduo Sang 已提交
1112 1113
            # self.assertAtMostOneSuccess(tasks, DropFixedTableTask) # not
            # really in case of drop-create-drop
1114 1115

        # if self._state.canAddData():
S
Shuduo Sang 已提交
1116 1117
        # self.assertIfExistThenSuccess(tasks, AddFixedDataTask)  # not true
        # actually
1118 1119 1120 1121 1122 1123

        # if self._state.canReadData():
            # Nothing for sure

        newState = self._findCurrentState()
        logger.debug("[STT] New DB state determined: {}".format(newState))
S
Shuduo Sang 已提交
1124 1125
        # can old state move to new state through the tasks?
        self._curState.verifyTasksToState(tasks, newState)
1126 1127 1128
        self._curState = newState

    def pickTaskType(self):
S
Shuduo Sang 已提交
1129 1130
        # all the task types we can choose from at curent state
        taskTypes = self.getTaskTypes()
1131 1132 1133
        weights = []
        for tt in taskTypes:
            endState = tt.getEndState()
S
Shuduo Sang 已提交
1134 1135 1136
            if endState is not None:
                # TODO: change to a method
                weights.append(self._stateWeights[endState.getValIndex()])
1137
            else:
S
Shuduo Sang 已提交
1138 1139
                # read data task, default to 10: TODO: change to a constant
                weights.append(10)
1140
        i = self._weighted_choice_sub(weights)
S
Shuduo Sang 已提交
1141
        # logger.debug(" (weighted random:{}/{}) ".format(i, len(taskTypes)))
1142 1143
        return taskTypes[i]

S
Shuduo Sang 已提交
1144 1145 1146 1147 1148
    # ref:
    # https://eli.thegreenplace.net/2010/01/22/weighted-random-generation-in-python/
    def _weighted_choice_sub(self, weights):
        # TODO: use our dice to ensure it being determinstic?
        rnd = random.random() * sum(weights)
1149 1150 1151 1152
        for i, w in enumerate(weights):
            rnd -= w
            if rnd < 0:
                return i
1153

1154
# Manager of the Database Data/Connection
S
Shuduo Sang 已提交
1155 1156 1157 1158


class DbManager():
    def __init__(self, resetDb=True):
S
Steven Li 已提交
1159
        self.tableNumQueue = LinearQueue()
S
Shuduo Sang 已提交
1160 1161 1162
        # datetime.datetime(2019, 1, 1) # initial date time tick
        self._lastTick = self.setupLastTick()
        self._lastInt = 0  # next one is initial integer
1163
        self._lock = threading.RLock()
S
Shuduo Sang 已提交
1164

1165
        # self.openDbServerConnection()
S
Shuduo Sang 已提交
1166 1167
        self._dbConn = DbConn.createNative() if (
            gConfig.connector_type == 'native') else DbConn.createRest()
1168
        try:
S
Shuduo Sang 已提交
1169
            self._dbConn.open()  # may throw taos.error.ProgrammingError: disconnected
1170 1171
        except taos.error.ProgrammingError as err:
            # print("Error type: {}, msg: {}, value: {}".format(type(err), err.msg, err))
S
Shuduo Sang 已提交
1172 1173 1174
            if (err.msg == 'client disconnected'):  # cannot open DB connection
                print(
                    "Cannot establish DB connection, please re-run script without parameter, and follow the instructions.")
S
Steven Li 已提交
1175
                sys.exit(2)
1176
            else:
S
Shuduo Sang 已提交
1177 1178
                raise
        except BaseException:
S
Steven Li 已提交
1179
            print("[=] Unexpected exception")
S
Shuduo Sang 已提交
1180 1181 1182 1183
            raise

        if resetDb:
            self._dbConn.resetDb()  # drop and recreate DB
1184

S
Shuduo Sang 已提交
1185 1186
        # Do this after dbConn is in proper shape
        self._stateMachine = StateMechine(self._dbConn)
1187

1188 1189 1190
    def getDbConn(self):
        return self._dbConn

S
Shuduo Sang 已提交
1191
    def getStateMachine(self) -> StateMechine:
1192 1193 1194 1195
        return self._stateMachine

    # def getState(self):
    #     return self._stateMachine.getCurrentState()
1196 1197 1198 1199 1200 1201

    # We aim to create a starting time tick, such that, whenever we run our test here once
    # We should be able to safely create 100,000 records, which will not have any repeated time stamp
    # when we re-run the test in 3 minutes (180 seconds), basically we should expand time duration
    # by a factor of 500.
    # TODO: what if it goes beyond 10 years into the future
1202
    # TODO: fix the error as result of above: "tsdb timestamp is out of range"
1203
    def setupLastTick(self):
1204
        t1 = datetime.datetime(2020, 6, 1)
1205
        t2 = datetime.datetime.now()
S
Shuduo Sang 已提交
1206 1207 1208 1209
        # maybe a very large number, takes 69 years to exceed Python int range
        elSec = int(t2.timestamp() - t1.timestamp())
        elSec2 = (elSec % (8 * 12 * 30 * 24 * 60 * 60 / 500)) * \
            500  # a number representing seconds within 10 years
1210
        # print("elSec = {}".format(elSec))
S
Shuduo Sang 已提交
1211 1212 1213
        t3 = datetime.datetime(2012, 1, 1)  # default "keep" is 10 years
        t4 = datetime.datetime.fromtimestamp(
            t3.timestamp() + elSec2)  # see explanation above
1214 1215 1216
        logger.info("Setting up TICKS to start from: {}".format(t4))
        return t4

S
Shuduo Sang 已提交
1217
    def pickAndAllocateTable(self):  # pick any table, and "use" it
S
Steven Li 已提交
1218 1219
        return self.tableNumQueue.pickAndAllocate()

1220 1221 1222 1223 1224
    def addTable(self):
        with self._lock:
            tIndex = self.tableNumQueue.push()
        return tIndex

1225 1226
    @classmethod
    def getFixedSuperTableName(cls):
1227
        return "fs_table"
1228

S
Shuduo Sang 已提交
1229
    def releaseTable(self, i):  # return the table back, so others can use it
S
Steven Li 已提交
1230 1231
        self.tableNumQueue.release(i)

1232
    def getNextTick(self):
S
Shuduo Sang 已提交
1233 1234
        with self._lock:  # prevent duplicate tick
            if Dice.throw(10) == 0:  # 1 in 10 chance
S
Steven Li 已提交
1235
                return self._lastTick + datetime.timedelta(0, -100)
S
Shuduo Sang 已提交
1236 1237 1238
            else:  # regular
                # add one second to it
                self._lastTick += datetime.timedelta(0, 1)
S
Steven Li 已提交
1239
                return self._lastTick
1240 1241

    def getNextInt(self):
1242 1243 1244
        with self._lock:
            self._lastInt += 1
            return self._lastInt
1245 1246

    def getNextBinary(self):
S
Shuduo Sang 已提交
1247 1248
        return "Beijing_Shanghai_Los_Angeles_New_York_San_Francisco_Chicago_Beijing_Shanghai_Los_Angeles_New_York_San_Francisco_Chicago_{}".format(
            self.getNextInt())
1249 1250 1251

    def getNextFloat(self):
        return 0.9 + self.getNextInt()
S
Shuduo Sang 已提交
1252

S
Steven Li 已提交
1253
    def getTableNameToDelete(self):
S
Shuduo Sang 已提交
1254 1255
        tblNum = self.tableNumQueue.pop()  # TODO: race condition!
        if (not tblNum):  # maybe false
1256
            return False
S
Shuduo Sang 已提交
1257

S
Steven Li 已提交
1258 1259
        return "table_{}".format(tblNum)

1260
    def cleanUp(self):
S
Shuduo Sang 已提交
1261 1262
        self._dbConn.close()

1263

1264
class TaskExecutor():
1265
    class BoundedList:
S
Shuduo Sang 已提交
1266
        def __init__(self, size=10):
1267 1268 1269
            self._size = size
            self._list = []

S
Shuduo Sang 已提交
1270 1271
        def add(self, n: int):
            if not self._list:  # empty
1272 1273 1274 1275 1276 1277 1278
                self._list.append(n)
                return
            # now we should insert
            nItems = len(self._list)
            insPos = 0
            for i in range(nItems):
                insPos = i
S
Shuduo Sang 已提交
1279 1280 1281
                if n <= self._list[i]:  # smaller than this item, time to insert
                    break  # found the insertion point
                insPos += 1  # insert to the right
1282

S
Shuduo Sang 已提交
1283 1284
            if insPos == 0:  # except for the 1st item, # TODO: elimiate first item as gating item
                return  # do nothing
1285 1286

            # print("Inserting at postion {}, value: {}".format(insPos, n))
S
Shuduo Sang 已提交
1287 1288
            self._list.insert(insPos, n)  # insert

1289
            newLen = len(self._list)
S
Shuduo Sang 已提交
1290 1291 1292 1293 1294
            if newLen <= self._size:
                return  # do nothing
            elif newLen == (self._size + 1):
                del self._list[0]  # remove the first item
            else:
1295 1296 1297 1298 1299 1300 1301
                raise RuntimeError("Corrupt Bounded List")

        def __str__(self):
            return repr(self._list)

    _boundedList = BoundedList()

1302 1303 1304
    def __init__(self, curStep):
        self._curStep = curStep

1305 1306 1307 1308
    @classmethod
    def getBoundedList(cls):
        return cls._boundedList

1309 1310 1311
    def getCurStep(self):
        return self._curStep

S
Shuduo Sang 已提交
1312
    def execute(self, task: Task, wt: WorkerThread):  # execute a task on a thread
1313
        task.execute(wt)
1314

1315 1316 1317 1318
    def recordDataMark(self, n: int):
        # print("[{}]".format(n), end="", flush=True)
        self._boundedList.add(n)

1319 1320
    # def logInfo(self, msg):
    #     logger.info("    T[{}.x]: ".format(self._curStep) + msg)
1321

1322 1323
    # def logDebug(self, msg):
    #     logger.debug("    T[{}.x]: ".format(self._curStep) + msg)
1324

S
Shuduo Sang 已提交
1325

S
Steven Li 已提交
1326
class Task():
1327 1328 1329 1330
    taskSn = 100

    @classmethod
    def allocTaskNum(cls):
S
Shuduo Sang 已提交
1331
        Task.taskSn += 1  # IMPORTANT: cannot use cls.taskSn, since each sub class will have a copy
S
Steven Li 已提交
1332 1333
        # logger.debug("Allocating taskSN: {}".format(Task.taskSn))
        return Task.taskSn
1334

S
Shuduo Sang 已提交
1335
    def __init__(self, dbManager: DbManager, execStats: ExecutionStats):
1336
        self._dbManager = dbManager
S
Shuduo Sang 已提交
1337
        self._workerThread = None
1338
        self._err = None
1339
        self._aborted = False
1340
        self._curStep = None
S
Shuduo Sang 已提交
1341
        self._numRows = None  # Number of rows affected
1342

S
Shuduo Sang 已提交
1343
        # Assign an incremental task serial number
1344
        self._taskNum = self.allocTaskNum()
S
Steven Li 已提交
1345
        # logger.debug("Creating new task {}...".format(self._taskNum))
1346

1347
        self._execStats = execStats
S
Shuduo Sang 已提交
1348
        self._lastSql = ""  # last SQL executed/attempted
1349

1350
    def isSuccess(self):
S
Shuduo Sang 已提交
1351
        return self._err is None
1352

1353 1354 1355
    def isAborted(self):
        return self._aborted

S
Shuduo Sang 已提交
1356
    def clone(self):  # TODO: why do we need this again?
1357
        newTask = self.__class__(self._dbManager, self._execStats)
1358 1359 1360
        return newTask

    def logDebug(self, msg):
S
Shuduo Sang 已提交
1361 1362 1363
        self._workerThread.logDebug(
            "Step[{}.{}] {}".format(
                self._curStep, self._taskNum, msg))
1364 1365

    def logInfo(self, msg):
S
Shuduo Sang 已提交
1366 1367 1368
        self._workerThread.logInfo(
            "Step[{}.{}] {}".format(
                self._curStep, self._taskNum, msg))
1369

1370
    def _executeInternal(self, te: TaskExecutor, wt: WorkerThread):
S
Shuduo Sang 已提交
1371 1372 1373
        raise RuntimeError(
            "To be implemeted by child classes, class name: {}".format(
                self.__class__.__name__))
1374

1375 1376
    def execute(self, wt: WorkerThread):
        wt.verifyThreadSelf()
S
Shuduo Sang 已提交
1377
        self._workerThread = wt  # type: ignore
1378 1379

        te = wt.getTaskExecutor()
1380
        self._curStep = te.getCurStep()
S
Shuduo Sang 已提交
1381 1382
        self.logDebug(
            "[-] executing task {}...".format(self.__class__.__name__))
1383 1384

        self._err = None
S
Shuduo Sang 已提交
1385 1386
        self._execStats.beginTaskType(
            self.__class__.__name__)  # mark beginning
1387
        try:
S
Shuduo Sang 已提交
1388
            self._executeInternal(te, wt)  # TODO: no return value?
1389
        except taos.error.ProgrammingError as err:
S
Shuduo Sang 已提交
1390 1391 1392 1393 1394 1395 1396 1397
            errno2 = err.errno if (
                err.errno > 0) else 0x80000000 + err.errno  # correct error scheme
            if (errno2 in [0x200, 0x360, 0x362, 0x36A, 0x36B, 0x36D, 0x381, 0x380, 0x383, 0x503, 0x600,
                           1000  # REST catch-all error
                           ]):  # allowed errors
                self.logDebug(
                    "[=] Acceptable Taos library exception: errno=0x{:X}, msg: {}, SQL: {}".format(
                        errno2, err, self._lastSql))
1398
                print("_", end="", flush=True)
S
Shuduo Sang 已提交
1399
                self._err = err
1400
            else:
S
Shuduo Sang 已提交
1401 1402
                errMsg = "[=] Unexpected Taos library exception: errno=0x{:X}, msg: {}, SQL: {}".format(
                    errno2, err, self._lastSql)
1403
                self.logDebug(errMsg)
S
Shuduo Sang 已提交
1404 1405 1406 1407 1408
                if gConfig.debug:
                    raise  # so that we see full stack
                else:  # non-debug
                    print(
                        "\n\n----------------------------\nProgram ABORTED Due to Unexpected TAOS Error: \n\n{}\n".format(errMsg) +
1409
                        "----------------------------\n")
1410 1411 1412
                    # sys.exit(-1)
                    self._err = err
                    self._aborted = True
S
Shuduo Sang 已提交
1413
        except Exception as e:
S
Steven Li 已提交
1414
            self.logInfo("Non-TAOS exception encountered")
S
Shuduo Sang 已提交
1415
            self._err = e
S
Steven Li 已提交
1416 1417
            self._aborted = True
            traceback.print_exc()
S
Shuduo Sang 已提交
1418 1419 1420 1421
        except BaseException:
            self.logDebug(
                "[=] Unexpected exception, SQL: {}".format(
                    self._lastSql))
1422
            raise
1423
        self._execStats.endTaskType(self.__class__.__name__, self.isSuccess())
S
Shuduo Sang 已提交
1424 1425 1426 1427 1428

        self.logDebug("[X] task execution completed, {}, status: {}".format(
            self.__class__.__name__, "Success" if self.isSuccess() else "Failure"))
        # TODO: merge with above.
        self._execStats.incExecCount(self.__class__.__name__, self.isSuccess())
S
Steven Li 已提交
1429

1430
    def execSql(self, sql):
1431
        self._lastSql = sql
1432
        return self._dbManager.execute(sql)
1433

S
Shuduo Sang 已提交
1434
    def execWtSql(self, wt: WorkerThread, sql):  # execute an SQL on the worker thread
1435 1436 1437
        self._lastSql = sql
        return wt.execSql(sql)

S
Shuduo Sang 已提交
1438
    def queryWtSql(self, wt: WorkerThread, sql):  # execute an SQL on the worker thread
1439 1440 1441
        self._lastSql = sql
        return wt.querySql(sql)

S
Shuduo Sang 已提交
1442
    def getQueryResult(self, wt: WorkerThread):  # execute an SQL on the worker thread
1443 1444 1445
        return wt.getQueryResult()


1446
class ExecutionStats:
1447
    def __init__(self):
S
Shuduo Sang 已提交
1448 1449
        # total/success times for a task
        self._execTimes: Dict[str, [int, int]] = {}
1450 1451 1452
        self._tasksInProgress = 0
        self._lock = threading.Lock()
        self._firstTaskStartTime = None
1453
        self._execStartTime = None
S
Shuduo Sang 已提交
1454 1455
        self._elapsedTime = 0.0  # total elapsed time
        self._accRunTime = 0.0  # accumulated run time
1456

1457 1458 1459
        self._failed = False
        self._failureReason = None

S
Steven Li 已提交
1460
    def __str__(self):
S
Shuduo Sang 已提交
1461 1462
        return "[ExecStats: _failed={}, _failureReason={}".format(
            self._failed, self._failureReason)
S
Steven Li 已提交
1463 1464

    def isFailed(self):
S
Shuduo Sang 已提交
1465
        return self._failed
S
Steven Li 已提交
1466

1467 1468 1469 1470 1471 1472
    def startExec(self):
        self._execStartTime = time.time()

    def endExec(self):
        self._elapsedTime = time.time() - self._execStartTime

S
Shuduo Sang 已提交
1473
    def incExecCount(self, klassName, isSuccess):  # TODO: add a lock here
1474 1475
        if klassName not in self._execTimes:
            self._execTimes[klassName] = [0, 0]
S
Shuduo Sang 已提交
1476 1477
        t = self._execTimes[klassName]  # tuple for the data
        t[0] += 1  # index 0 has the "total" execution times
1478
        if isSuccess:
S
Shuduo Sang 已提交
1479
            t[1] += 1  # index 1 has the "success" execution times
1480 1481 1482

    def beginTaskType(self, klassName):
        with self._lock:
S
Shuduo Sang 已提交
1483 1484
            if self._tasksInProgress == 0:  # starting a new round
                self._firstTaskStartTime = time.time()  # I am now the first task
1485 1486 1487 1488 1489
            self._tasksInProgress += 1

    def endTaskType(self, klassName, isSuccess):
        with self._lock:
            self._tasksInProgress -= 1
S
Shuduo Sang 已提交
1490
            if self._tasksInProgress == 0:  # all tasks have stopped
1491 1492 1493
                self._accRunTime += (time.time() - self._firstTaskStartTime)
                self._firstTaskStartTime = None

1494 1495 1496 1497
    def registerFailure(self, reason):
        self._failed = True
        self._failureReason = reason

1498
    def printStats(self):
S
Shuduo Sang 已提交
1499 1500 1501 1502 1503 1504
        logger.info(
            "----------------------------------------------------------------------")
        logger.info(
            "| Crash_Gen test {}, with the following stats:". format(
                "FAILED (reason: {})".format(
                    self._failureReason) if self._failed else "SUCCEEDED"))
1505
        logger.info("| Task Execution Times (success/total):")
1506
        execTimesAny = 0
S
Shuduo Sang 已提交
1507
        for k, n in self._execTimes.items():
1508
            execTimesAny += n[0]
S
Shuduo Sang 已提交
1509 1510 1511 1512 1513 1514 1515 1516 1517 1518 1519 1520 1521 1522 1523 1524 1525 1526 1527 1528
            logger.info("|    {0:<24}: {1}/{2}".format(k, n[1], n[0]))

        logger.info(
            "| Total Tasks Executed (success or not): {} ".format(execTimesAny))
        logger.info(
            "| Total Tasks In Progress at End: {}".format(
                self._tasksInProgress))
        logger.info(
            "| Total Task Busy Time (elapsed time when any task is in progress): {:.3f} seconds".format(
                self._accRunTime))
        logger.info(
            "| Average Per-Task Execution Time: {:.3f} seconds".format(self._accRunTime / execTimesAny))
        logger.info(
            "| Total Elapsed Time (from wall clock): {:.3f} seconds".format(
                self._elapsedTime))
        logger.info(
            "| Top numbers written: {}".format(
                TaskExecutor.getBoundedList()))
        logger.info(
            "----------------------------------------------------------------------")
1529 1530 1531


class StateTransitionTask(Task):
1532 1533 1534 1535 1536
    LARGE_NUMBER_OF_TABLES = 35
    SMALL_NUMBER_OF_TABLES = 3
    LARGE_NUMBER_OF_RECORDS = 50
    SMALL_NUMBER_OF_RECORDS = 3

1537
    @classmethod
S
Shuduo Sang 已提交
1538
    def getInfo(cls):  # each sub class should supply their own information
1539 1540
        raise RuntimeError("Overriding method expected")

S
Shuduo Sang 已提交
1541
    _endState = None
1542
    @classmethod
S
Shuduo Sang 已提交
1543
    def getEndState(cls):  # TODO: optimize by calling it fewer times
1544 1545
        raise RuntimeError("Overriding method expected")

1546 1547 1548
    # @classmethod
    # def getBeginStates(cls):
    #     return cls.getInfo()[0]
1549

1550 1551 1552
    # @classmethod
    # def getEndState(cls): # returning the class name
    #     return cls.getInfo()[0]
1553 1554

    @classmethod
1555 1556 1557
    def canBeginFrom(cls, state: AnyState):
        # return state.getValue() in cls.getBeginStates()
        raise RuntimeError("must be overriden")
1558

1559 1560 1561 1562
    @classmethod
    def getRegTableName(cls, i):
        return "db.reg_table_{}".format(i)

1563 1564
    def execute(self, wt: WorkerThread):
        super().execute(wt)
S
Shuduo Sang 已提交
1565 1566


1567
class TaskCreateDb(StateTransitionTask):
1568
    @classmethod
1569
    def getEndState(cls):
S
Shuduo Sang 已提交
1570
        return StateDbOnly()
1571

1572 1573 1574 1575
    @classmethod
    def canBeginFrom(cls, state: AnyState):
        return state.canCreateDb()

1576
    def _executeInternal(self, te: TaskExecutor, wt: WorkerThread):
S
Shuduo Sang 已提交
1577 1578
        self.execWtSql(wt, "create database db")

1579

1580
class TaskDropDb(StateTransitionTask):
1581
    @classmethod
1582 1583
    def getEndState(cls):
        return StateEmpty()
1584

1585 1586 1587 1588
    @classmethod
    def canBeginFrom(cls, state: AnyState):
        return state.canDropDb()

1589
    def _executeInternal(self, te: TaskExecutor, wt: WorkerThread):
1590
        self.execWtSql(wt, "drop database db")
S
Steven Li 已提交
1591
        logger.debug("[OPS] database dropped at {}".format(time.time()))
1592

S
Shuduo Sang 已提交
1593

1594
class TaskCreateSuperTable(StateTransitionTask):
1595
    @classmethod
1596 1597
    def getEndState(cls):
        return StateSuperTableOnly()
1598

1599 1600
    @classmethod
    def canBeginFrom(cls, state: AnyState):
1601
        return state.canCreateFixedSuperTable()
1602

1603
    def _executeInternal(self, te: TaskExecutor, wt: WorkerThread):
S
Shuduo Sang 已提交
1604
        if not wt.dbInUse():  # no DB yet, to the best of our knowledge
1605 1606 1607
            logger.debug("Skipping task, no DB yet")
            return

S
Shuduo Sang 已提交
1608
        tblName = self._dbManager.getFixedSuperTableName()
1609
        # wt.execSql("use db")    # should always be in place
S
Shuduo Sang 已提交
1610 1611 1612 1613 1614
        self.execWtSql(
            wt,
            "create table db.{} (ts timestamp, speed int) tags (b binary(200), f float) ".format(tblName))
        # No need to create the regular tables, INSERT will do that
        # automatically
1615

S
Steven Li 已提交
1616

1617
class TaskReadData(StateTransitionTask):
1618
    @classmethod
1619
    def getEndState(cls):
S
Shuduo Sang 已提交
1620
        return None  # meaning doesn't affect state
1621

1622 1623 1624 1625
    @classmethod
    def canBeginFrom(cls, state: AnyState):
        return state.canReadData()

1626
    def _executeInternal(self, te: TaskExecutor, wt: WorkerThread):
S
Shuduo Sang 已提交
1627 1628 1629
        sTbName = self._dbManager.getFixedSuperTableName()
        self.queryWtSql(wt, "select TBNAME from db.{}".format(
            sTbName))  # TODO: analyze result set later
1630

S
Shuduo Sang 已提交
1631 1632
        if random.randrange(
                5) == 0:  # 1 in 5 chance, simulate a broken connection. TODO: break connection in all situations
1633 1634
            wt.getDbConn().close()
            wt.getDbConn().open()
1635
        else:
S
Shuduo Sang 已提交
1636 1637
            # wt.getDbConn().getQueryResult()
            rTables = self.getQueryResult(wt)
1638
            # print("rTables[0] = {}, type = {}".format(rTables[0], type(rTables[0])))
S
Shuduo Sang 已提交
1639
            for rTbName in rTables:  # regular tables
1640
                self.execWtSql(wt, "select * from db.{}".format(rTbName[0]))
1641

1642 1643
        # tdSql.query(" cars where tbname in ('carzero', 'carone')")

S
Shuduo Sang 已提交
1644

1645
class TaskDropSuperTable(StateTransitionTask):
1646
    @classmethod
1647
    def getEndState(cls):
S
Shuduo Sang 已提交
1648
        return StateDbOnly()
1649

1650 1651
    @classmethod
    def canBeginFrom(cls, state: AnyState):
1652
        return state.canDropFixedSuperTable()
1653

1654
    def _executeInternal(self, te: TaskExecutor, wt: WorkerThread):
S
Shuduo Sang 已提交
1655 1656 1657 1658 1659 1660 1661
        # 1/2 chance, we'll drop the regular tables one by one, in a randomized
        # sequence
        if Dice.throw(2) == 0:
            tblSeq = list(range(
                2 + (self.LARGE_NUMBER_OF_TABLES if gConfig.larger_data else self.SMALL_NUMBER_OF_TABLES)))
            random.shuffle(tblSeq)
            tickOutput = False  # if we have spitted out a "d" character for "drop regular table"
1662
            isSuccess = True
S
Shuduo Sang 已提交
1663 1664 1665
            for i in tblSeq:
                regTableName = self.getRegTableName(
                    i)  # "db.reg_table_{}".format(i)
1666
                try:
S
Shuduo Sang 已提交
1667 1668 1669 1670 1671 1672 1673
                    self.execWtSql(wt, "drop table {}".format(
                        regTableName))  # nRows always 0, like MySQL
                except taos.error.ProgrammingError as err:
                    # correcting for strange error number scheme
                    errno2 = err.errno if (
                        err.errno > 0) else 0x80000000 + err.errno
                    if (errno2 in [0x362]):  # mnode invalid table name
1674
                        isSuccess = False
S
Shuduo Sang 已提交
1675 1676 1677
                        logger.debug(
                            "[DB] Acceptable error when dropping a table")
                    continue  # try to delete next regular table
1678 1679

                if (not tickOutput):
S
Shuduo Sang 已提交
1680 1681
                    tickOutput = True  # Print only one time
                    if isSuccess:
1682 1683
                        print("d", end="", flush=True)
                    else:
S
Shuduo Sang 已提交
1684
                        print("f", end="", flush=True)
1685 1686

        # Drop the super table itself
S
Shuduo Sang 已提交
1687
        tblName = self._dbManager.getFixedSuperTableName()
1688
        self.execWtSql(wt, "drop table db.{}".format(tblName))
1689

S
Shuduo Sang 已提交
1690

1691 1692 1693
class TaskAlterTags(StateTransitionTask):
    @classmethod
    def getEndState(cls):
S
Shuduo Sang 已提交
1694
        return None  # meaning doesn't affect state
1695 1696 1697

    @classmethod
    def canBeginFrom(cls, state: AnyState):
S
Shuduo Sang 已提交
1698
        return state.canDropFixedSuperTable()  # if we can drop it, we can alter tags
1699 1700

    def _executeInternal(self, te: TaskExecutor, wt: WorkerThread):
S
Shuduo Sang 已提交
1701
        tblName = self._dbManager.getFixedSuperTableName()
1702
        dice = Dice.throw(4)
S
Shuduo Sang 已提交
1703
        if dice == 0:
1704
            sql = "alter table db.{} add tag extraTag int".format(tblName)
S
Shuduo Sang 已提交
1705
        elif dice == 1:
1706
            sql = "alter table db.{} drop tag extraTag".format(tblName)
S
Shuduo Sang 已提交
1707
        elif dice == 2:
1708
            sql = "alter table db.{} drop tag newTag".format(tblName)
S
Shuduo Sang 已提交
1709 1710 1711
        else:  # dice == 3
            sql = "alter table db.{} change tag extraTag newTag".format(
                tblName)
1712 1713

        self.execWtSql(wt, sql)
1714

S
Shuduo Sang 已提交
1715

1716
class TaskAddData(StateTransitionTask):
S
Shuduo Sang 已提交
1717 1718
    # Track which table is being actively worked on
    activeTable: Set[int] = set()
1719

S
Shuduo Sang 已提交
1720 1721
    # We use these two files to record operations to DB, useful for power-off
    # tests
1722 1723 1724 1725 1726
    fAddLogReady = None
    fAddLogDone = None

    @classmethod
    def prepToRecordOps(cls):
S
Shuduo Sang 已提交
1727 1728 1729 1730
        if gConfig.record_ops:
            if (cls.fAddLogReady is None):
                logger.info(
                    "Recording in a file operations to be performed...")
1731
                cls.fAddLogReady = open("add_log_ready.txt", "w")
S
Shuduo Sang 已提交
1732
            if (cls.fAddLogDone is None):
1733 1734
                logger.info("Recording in a file operations completed...")
                cls.fAddLogDone = open("add_log_done.txt", "w")
1735

1736
    @classmethod
1737 1738
    def getEndState(cls):
        return StateHasData()
1739 1740 1741 1742

    @classmethod
    def canBeginFrom(cls, state: AnyState):
        return state.canAddData()
S
Shuduo Sang 已提交
1743

1744
    def _executeInternal(self, te: TaskExecutor, wt: WorkerThread):
1745
        ds = self._dbManager
S
Shuduo Sang 已提交
1746 1747 1748 1749 1750 1751 1752 1753 1754 1755 1756
        # wt.execSql("use db") # TODO: seems to be an INSERT bug to require
        # this
        tblSeq = list(
            range(
                self.LARGE_NUMBER_OF_TABLES if gConfig.larger_data else self.SMALL_NUMBER_OF_TABLES))
        random.shuffle(tblSeq)
        for i in tblSeq:
            if (i in self.activeTable):  # wow already active
                # logger.info("Concurrent data insertion into table: {}".format(i))
                # print("ct({})".format(i), end="", flush=True) # Concurrent
                # insertion into table
1757 1758
                print("x", end="", flush=True)
            else:
S
Shuduo Sang 已提交
1759 1760 1761 1762 1763 1764 1765 1766
                self.activeTable.add(i)  # marking it active
            # No need to shuffle data sequence, unless later we decide to do
            # non-increment insertion
            regTableName = self.getRegTableName(
                i)  # "db.reg_table_{}".format(i)
            for j in range(
                    self.LARGE_NUMBER_OF_RECORDS if gConfig.larger_data else self.SMALL_NUMBER_OF_RECORDS):  # number of records per table
                nextInt = ds.getNextInt()
1767 1768
                if gConfig.record_ops:
                    self.prepToRecordOps()
S
Shuduo Sang 已提交
1769 1770 1771
                    self.fAddLogReady.write(
                        "Ready to write {} to {}\n".format(
                            nextInt, regTableName))
1772 1773 1774
                    self.fAddLogReady.flush()
                    os.fsync(self.fAddLogReady)
                sql = "insert into {} using {} tags ('{}', {}) values ('{}', {});".format(
S
Shuduo Sang 已提交
1775 1776
                    regTableName,
                    ds.getFixedSuperTableName(),
1777
                    ds.getNextBinary(), ds.getNextFloat(),
1778
                    ds.getNextTick(), nextInt)
S
Shuduo Sang 已提交
1779 1780 1781
                self.execWtSql(wt, sql)
                # Successfully wrote the data into the DB, let's record it
                # somehow
1782
                te.recordDataMark(nextInt)
1783
                if gConfig.record_ops:
S
Shuduo Sang 已提交
1784 1785 1786
                    self.fAddLogDone.write(
                        "Wrote {} to {}\n".format(
                            nextInt, regTableName))
1787 1788
                    self.fAddLogDone.flush()
                    os.fsync(self.fAddLogDone)
S
Shuduo Sang 已提交
1789
            self.activeTable.discard(i)  # not raising an error, unlike remove
1790 1791


S
Steven Li 已提交
1792 1793
# Deterministic random number generator
class Dice():
S
Shuduo Sang 已提交
1794
    seeded = False  # static, uninitialized
S
Steven Li 已提交
1795 1796

    @classmethod
S
Shuduo Sang 已提交
1797
    def seed(cls, s):  # static
S
Steven Li 已提交
1798
        if (cls.seeded):
S
Shuduo Sang 已提交
1799 1800
            raise RuntimeError(
                "Cannot seed the random generator more than once")
S
Steven Li 已提交
1801 1802 1803 1804 1805
        cls.verifyRNG()
        random.seed(s)
        cls.seeded = True  # TODO: protect against multi-threading

    @classmethod
S
Shuduo Sang 已提交
1806
    def verifyRNG(cls):  # Verify that the RNG is determinstic
S
Steven Li 已提交
1807 1808 1809 1810
        random.seed(0)
        x1 = random.randrange(0, 1000)
        x2 = random.randrange(0, 1000)
        x3 = random.randrange(0, 1000)
S
Shuduo Sang 已提交
1811
        if (x1 != 864 or x2 != 394 or x3 != 776):
S
Steven Li 已提交
1812 1813 1814
            raise RuntimeError("System RNG is not deterministic")

    @classmethod
S
Shuduo Sang 已提交
1815
    def throw(cls, stop):  # get 0 to stop-1
1816
        return cls.throwRange(0, stop)
S
Steven Li 已提交
1817 1818

    @classmethod
S
Shuduo Sang 已提交
1819 1820
    def throwRange(cls, start, stop):  # up to stop-1
        if (not cls.seeded):
S
Steven Li 已提交
1821
            raise RuntimeError("Cannot throw dice before seeding it")
1822
        return random.randrange(start, stop)
S
Steven Li 已提交
1823 1824 1825


# Anyone needing to carry out work should simply come here
1826 1827 1828 1829 1830 1831 1832 1833 1834 1835
# class WorkDispatcher():
#     def __init__(self, dbState):
#         # self.totalNumMethods = 2
#         self.tasks = [
#             # CreateTableTask(dbState), # Obsolete
#             # DropTableTask(dbState),
#             # AddDataTask(dbState),
#         ]

#     def throwDice(self):
S
Shuduo Sang 已提交
1836
#         max = len(self.tasks) - 1
1837 1838 1839 1840 1841 1842 1843 1844 1845 1846 1847
#         dRes = random.randint(0, max)
#         # logger.debug("Threw the dice in range [{},{}], and got: {}".format(0,max,dRes))
#         return dRes

#     def pickTask(self):
#         dice = self.throwDice()
#         return self.tasks[dice]

#     def doWork(self, workerThread):
#         task = self.pickTask()
#         task.execute(workerThread)
S
Steven Li 已提交
1848

S
Steven Li 已提交
1849 1850
class LoggingFilter(logging.Filter):
    def filter(self, record: logging.LogRecord):
S
Shuduo Sang 已提交
1851 1852
        if (record.levelno >= logging.INFO):
            return True  # info or above always log
S
Steven Li 已提交
1853

S
Steven Li 已提交
1854 1855 1856 1857
        # Commenting out below to adjust...

        # if msg.startswith("[TRD]"):
        #     return False
S
Steven Li 已提交
1858 1859
        return True

S
Shuduo Sang 已提交
1860 1861

class MyLoggingAdapter(logging.LoggerAdapter):
1862 1863 1864 1865
    def process(self, msg, kwargs):
        return "[{}]{}".format(threading.get_ident() % 10000, msg), kwargs
        # return '[%s] %s' % (self.extra['connid'], msg), kwargs

S
Shuduo Sang 已提交
1866 1867 1868

class SvcManager:

1869 1870 1871 1872 1873 1874 1875 1876 1877 1878 1879
    def __init__(self):
        print("Starting service manager")
        signal.signal(signal.SIGTERM, self.sigIntHandler)
        signal.signal(signal.SIGINT, self.sigIntHandler)
        self.ioThread = None
        self.subProcess = None
        self.shouldStop = False
        self.status = MainExec.STATUS_RUNNING

    def svcOutputReader(self, out: IO, queue):
        # print("This is the svcOutput Reader...")
S
Shuduo Sang 已提交
1880
        for line in out:  # iter(out.readline, b''):
1881
            # print("Finished reading a line: {}".format(line))
S
Shuduo Sang 已提交
1882 1883 1884
            queue.put(line.rstrip())  # get rid of new lines
        # meaning sub process must have died
        print("No more output from incoming IO")
1885 1886 1887
        out.close()

    def sigIntHandler(self, signalNumber, frame):
S
Shuduo Sang 已提交
1888
        if self.status != MainExec.STATUS_RUNNING:
1889
            print("Ignoring repeated SIGINT...")
S
Shuduo Sang 已提交
1890 1891
            return  # do nothing if it's already not running
        self.status = MainExec.STATUS_STOPPING  # immediately set our status
1892 1893 1894 1895 1896 1897 1898

        print("Terminating program...")
        self.subProcess.send_signal(signal.SIGINT)
        self.shouldStop = True
        self.joinIoThread()

    def joinIoThread(self):
S
Shuduo Sang 已提交
1899
        if self.ioThread:
1900
            self.ioThread.join()
S
Shuduo Sang 已提交
1901
            self.ioThread = None
1902 1903 1904

    def run(self):
        ON_POSIX = 'posix' in sys.builtin_module_names
1905 1906 1907 1908 1909
        taosdPath = self.getBuildPath() + "/build/bin/taosd"
        cfgPath = self.getBuildPath() + "/test/cfg"

        print ("CBD: taosdPath:%s cfgPath:%s" % (taosPat, cfgPath))

1910 1911
        svcCmd = ['../../build/build/bin/taosd', '-c', '../../build/test/cfg']
        # svcCmd = ['vmstat', '1']
S
Shuduo Sang 已提交
1912 1913 1914 1915 1916 1917
        self.subProcess = subprocess.Popen(
            svcCmd,
            stdout=subprocess.PIPE,
            bufsize=1,
            close_fds=ON_POSIX,
            text=True)
1918
        q = Queue()
S
Shuduo Sang 已提交
1919 1920 1921 1922
        self.ioThread = threading.Thread(
            target=self.svcOutputReader, args=(
                self.subProcess.stdout, q))
        self.ioThread.daemon = True  # thread dies with the program
1923 1924
        self.ioThread.start()

S
Shuduo Sang 已提交
1925
        # proc = subprocess.Popen(['echo', '"to stdout"'],
1926 1927 1928 1929 1930
        #                 stdout=subprocess.PIPE,
        #                 )
        # stdout_value = proc.communicate()[0]
        # print('\tstdout: {}'.format(repr(stdout_value)))

S
Shuduo Sang 已提交
1931 1932 1933
        while True:
            try:
                line = q.get_nowait()  # getting output at fast speed
1934 1935
            except Empty:
                # print('no output yet')
S
Shuduo Sang 已提交
1936 1937
                time.sleep(2.3)  # wait only if there's no output
            else:  # got line
1938 1939 1940 1941 1942 1943 1944
                print(line)
            # print("----end of iteration----")
            if self.shouldStop:
                print("Ending main Svc thread")
                break

        print("end of loop")
S
Shuduo Sang 已提交
1945

1946 1947 1948
        self.joinIoThread()
        print("Finished")

S
Shuduo Sang 已提交
1949

1950 1951 1952 1953 1954 1955 1956 1957 1958 1959
class ClientManager:
    def __init__(self):
        print("Starting service manager")
        signal.signal(signal.SIGTERM, self.sigIntHandler)
        signal.signal(signal.SIGINT, self.sigIntHandler)

        self.status = MainExec.STATUS_RUNNING
        self.tc = None

    def sigIntHandler(self, signalNumber, frame):
S
Shuduo Sang 已提交
1960
        if self.status != MainExec.STATUS_RUNNING:
1961
            print("Ignoring repeated SIGINT...")
S
Shuduo Sang 已提交
1962 1963
            return  # do nothing if it's already not running
        self.status = MainExec.STATUS_STOPPING  # immediately set our status
1964 1965 1966 1967

        print("Terminating program...")
        self.tc.requestToStop()

S
Shuduo Sang 已提交
1968
    def _printLastNumbers(self):  # to verify data durability
1969 1970
        dbManager = DbManager(resetDb=False)
        dbc = dbManager.getDbConn()
S
Shuduo Sang 已提交
1971
        if dbc.query("show databases") == 0:  # no databae
1972
            return
S
Shuduo Sang 已提交
1973
        if dbc.query("show tables") == 0:  # no tables
S
Steven Li 已提交
1974
            return
1975 1976

        dbc.execute("use db")
S
Shuduo Sang 已提交
1977
        sTbName = dbManager.getFixedSuperTableName()
1978 1979

        # get all regular tables
S
Shuduo Sang 已提交
1980 1981
        # TODO: analyze result set later
        dbc.query("select TBNAME from db.{}".format(sTbName))
1982 1983 1984
        rTables = dbc.getQueryResult()

        bList = TaskExecutor.BoundedList()
S
Shuduo Sang 已提交
1985
        for rTbName in rTables:  # regular tables
1986 1987
            dbc.query("select speed from db.{}".format(rTbName[0]))
            numbers = dbc.getQueryResult()
S
Shuduo Sang 已提交
1988
            for row in numbers:
1989 1990 1991 1992 1993 1994
                # print("<{}>".format(n), end="", flush=True)
                bList.add(row[0])

        print("Top numbers in DB right now: {}".format(bList))
        print("TDengine client execution is about to start in 2 seconds...")
        time.sleep(2.0)
S
Shuduo Sang 已提交
1995
        dbManager = None  # release?
1996 1997 1998 1999 2000 2001 2002

    def prepare(self):
        self._printLastNumbers()

    def run(self):
        self._printLastNumbers()

S
Shuduo Sang 已提交
2003 2004
        dbManager = DbManager()  # Regular function
        Dice.seed(0)  # initial seeding of dice
2005
        thPool = ThreadPool(gConfig.num_threads, gConfig.max_steps)
2006
        self.tc = ThreadCoordinator(thPool, dbManager)
S
Shuduo Sang 已提交
2007

2008
        self.tc.run()
S
Steven Li 已提交
2009 2010
        # print("exec stats: {}".format(self.tc.getExecStats()))
        # print("TC failed = {}".format(self.tc.isFailed()))
2011
        self.conclude()
S
Steven Li 已提交
2012
        # print("TC failed (2) = {}".format(self.tc.isFailed()))
S
Shuduo Sang 已提交
2013 2014
        # Linux return code: ref https://shapeshed.com/unix-exit-codes/
        return 1 if self.tc.isFailed() else 0
2015 2016 2017

    def conclude(self):
        self.tc.printStats()
S
Shuduo Sang 已提交
2018
        self.tc.getDbManager().cleanUp()
2019 2020 2021 2022 2023 2024 2025 2026 2027 2028


class MainExec:
    STATUS_RUNNING = 1
    STATUS_STOPPING = 2
    # STATUS_STOPPED = 3 # Not used yet

    @classmethod
    def runClient(cls):
        clientManager = ClientManager()
S
Steven Li 已提交
2029
        return clientManager.run()
2030 2031 2032

    @classmethod
    def runService(cls):
2033 2034
        svcManager = SvcManager()
        svcManager.run()
2035 2036

    @classmethod
S
Shuduo Sang 已提交
2037
    def runTemp(cls):  # for debugging purposes
2038 2039
        # # Hack to exercise reading from disk, imcreasing coverage. TODO: fix
        # dbc = dbState.getDbConn()
S
Shuduo Sang 已提交
2040
        # sTbName = dbState.getFixedSuperTableName()
2041 2042
        # dbc.execute("create database if not exists db")
        # if not dbState.getState().equals(StateEmpty()):
S
Shuduo Sang 已提交
2043
        #     dbc.execute("use db")
2044 2045 2046 2047 2048 2049 2050 2051 2052 2053 2054

        # rTables = None
        # try: # the super table may not exist
        #     sql = "select TBNAME from db.{}".format(sTbName)
        #     logger.info("Finding out tables in super table: {}".format(sql))
        #     dbc.query(sql) # TODO: analyze result set later
        #     logger.info("Fetching result")
        #     rTables = dbc.getQueryResult()
        #     logger.info("Result: {}".format(rTables))
        # except taos.error.ProgrammingError as err:
        #     logger.info("Initial Super table OPS error: {}".format(err))
S
Shuduo Sang 已提交
2055

2056 2057 2058 2059 2060 2061 2062 2063
        # # sys.exit()
        # if ( not rTables == None):
        #     # print("rTables[0] = {}, type = {}".format(rTables[0], type(rTables[0])))
        #     try:
        #         for rTbName in rTables : # regular tables
        #             ds = dbState
        #             logger.info("Inserting into table: {}".format(rTbName[0]))
        #             sql = "insert into db.{} values ('{}', {});".format(
S
Shuduo Sang 已提交
2064
        #                 rTbName[0],
2065 2066
        #                 ds.getNextTick(), ds.getNextInt())
        #             dbc.execute(sql)
S
Shuduo Sang 已提交
2067
        #         for rTbName in rTables : # regular tables
2068
        #             dbc.query("select * from db.{}".format(rTbName[0])) # TODO: check success failure
S
Shuduo Sang 已提交
2069
        #         logger.info("Initial READING operation is successful")
2070
        #     except taos.error.ProgrammingError as err:
S
Shuduo Sang 已提交
2071 2072
        #         logger.info("Initial WRITE/READ error: {}".format(err))

2073 2074 2075
        # Sandbox testing code
        # dbc = dbState.getDbConn()
        # while True:
S
Shuduo Sang 已提交
2076
        #     rows = dbc.query("show databases")
2077
        #     print("Rows: {}, time={}".format(rows, time.time()))
S
Shuduo Sang 已提交
2078 2079
        return

S
Steven Li 已提交
2080

2081
def main():
S
Shuduo Sang 已提交
2082 2083
    # Super cool Python argument library:
    # https://docs.python.org/3/library/argparse.html
2084 2085 2086 2087 2088 2089 2090 2091 2092
    parser = argparse.ArgumentParser(
        formatter_class=argparse.RawDescriptionHelpFormatter,
        description=textwrap.dedent('''\
            TDengine Auto Crash Generator (PLEASE NOTICE the Prerequisites Below)
            ---------------------------------------------------------------------
            1. You build TDengine in the top level ./build directory, as described in offical docs
            2. You run the server there before this script: ./build/bin/taosd -c test/cfg

            '''))
2093

S
Shuduo Sang 已提交
2094 2095 2096 2097 2098 2099 2100 2101 2102 2103 2104 2105 2106 2107 2108 2109 2110 2111 2112 2113 2114 2115 2116 2117 2118 2119 2120 2121 2122 2123 2124 2125 2126 2127 2128 2129 2130 2131 2132 2133 2134 2135 2136 2137 2138 2139
    parser.add_argument(
        '-c',
        '--connector-type',
        action='store',
        default='native',
        type=str,
        help='Connector type to use: native, rest, or mixed (default: 10)')
    parser.add_argument(
        '-d',
        '--debug',
        action='store_true',
        help='Turn on DEBUG mode for more logging (default: false)')
    parser.add_argument(
        '-e',
        '--run-tdengine',
        action='store_true',
        help='Run TDengine service in foreground (default: false)')
    parser.add_argument(
        '-l',
        '--larger-data',
        action='store_true',
        help='Write larger amount of data during write operations (default: false)')
    parser.add_argument(
        '-p',
        '--per-thread-db-connection',
        action='store_true',
        help='Use a single shared db connection (default: false)')
    parser.add_argument(
        '-r',
        '--record-ops',
        action='store_true',
        help='Use a pair of always-fsynced fils to record operations performing + performed, for power-off tests (default: false)')
    parser.add_argument(
        '-s',
        '--max-steps',
        action='store',
        default=1000,
        type=int,
        help='Maximum number of steps to run (default: 100)')
    parser.add_argument(
        '-t',
        '--num-threads',
        action='store',
        default=5,
        type=int,
        help='Number of threads to run (default: 10)')
2140

2141
    global gConfig
2142
    gConfig = parser.parse_args()
2143

2144 2145 2146
    # if len(sys.argv) == 1:
    #     parser.print_help()
    #     sys.exit()
S
Shuduo Sang 已提交
2147

2148
    # Logging Stuff
2149
    global logger
S
Shuduo Sang 已提交
2150 2151
    _logger = logging.getLogger('CrashGen')  # real logger
    _logger.addFilter(LoggingFilter())
2152 2153 2154
    ch = logging.StreamHandler()
    _logger.addHandler(ch)

S
Shuduo Sang 已提交
2155 2156
    # Logging adapter, to be used as a logger
    logger = MyLoggingAdapter(_logger, [])
2157

S
Shuduo Sang 已提交
2158 2159
    if (gConfig.debug):
        logger.setLevel(logging.DEBUG)  # default seems to be INFO
S
Steven Li 已提交
2160 2161
    else:
        logger.setLevel(logging.INFO)
S
Shuduo Sang 已提交
2162

2163
    # Run server or client
S
Shuduo Sang 已提交
2164
    if gConfig.run_tdengine:  # run server
2165
        MainExec.runService()
S
Shuduo Sang 已提交
2166
    else:
S
Steven Li 已提交
2167
        return MainExec.runClient()
2168

S
Steven Li 已提交
2169
    # logger.info("Crash_Gen execution finished")
2170

S
Shuduo Sang 已提交
2171

2172
if __name__ == "__main__":
S
Steven Li 已提交
2173 2174 2175
    exitCode = main()
    # print("Exiting with code: {}".format(exitCode))
    sys.exit(exitCode)