From 12ccd1e586c3319fc31923a62b654bb423d76300 Mon Sep 17 00:00:00 2001 From: brb Date: Wed, 12 Jul 2023 13:59:53 +0200 Subject: [PATCH] added extra data flushing event with separate trigger time --- code/main.py | 119 +++++++++++++++++++++++++++++++++------------------ 1 file changed, 77 insertions(+), 42 deletions(-) diff --git a/code/main.py b/code/main.py index 5b84253..0ef9e24 100644 --- a/code/main.py +++ b/code/main.py @@ -14,45 +14,28 @@ from database import connect config = ConfigParser() config.read(Path(__file__).parent / 'defaults.ini') REPLICATES = int(config.get('main', 'replicates')) -USER_ID = int(config.get('main', 'user_id')) TRIGGER_INTERVAL_SECONDS = float( config.get( 'main', 'trigger_interval_seconds' ) ) +DATA_QUEUE_FLUSH_INTERVAL_SECONDS = float( + config.get( + 'main', + 'data_queue_flush_interval_seconds' + ) +) -def event( - replicates: int = REPLICATES, - user_id: int = USER_ID -): - logging.debug('started event') - # load environment variables - replicates = int(os.getenv('REPLICATES', default=REPLICATES)) - user_id = int(os.getenv('USER_ID', default=USER_ID)) - # setup data queue - data_queue: Queue = Queue() - # run event - for rep_num in range(replicates): - logging.debug(f'running replicate {rep_num+1} of {replicates}') - # do measurement - try: - data = measure_bandwidth() - except Exception as e: - logging.error(f'failed to measure bandwidth: {e}') - continue - # add data to queue - data_queue.put(data) - # polite pause - time.sleep(1) - # check if queue should be flushed - queue_size = data_queue.qsize() - if queue_size == 0: - logging.debug('queue empty, continuing') - return - # flush queue to database +def flush_data( + data_queue: Queue, + user_id: int +) -> None: + # connect to database db = connect() + # flush queue + queue_size = data_queue.qsize() for i in range(queue_size): # get data data = data_queue.get() @@ -78,26 +61,78 @@ def event( logging.debug('finished') +def event( + data_queue: Queue, + user_id: int, + replicates: int = REPLICATES +) -> None: + # run event + for rep_num in range(replicates): + logging.debug(f'running replicate {rep_num+1} of {replicates}') + # do measurement + try: + data = measure_bandwidth() + except Exception as e: + logging.error(f'failed to measure bandwidth: {e}') + continue + # add data to queue + data_queue.put(data) + # polite pause + time.sleep(1) + # flush queue + if data_queue.qsize() > 0: + flush_data( + data_queue=data_queue, + user_id=user_id + ) + logging.debug('finished') + + if __name__ == '__main__': + # setup initialise_app() + user_id = int(os.getenv('USER_ID', default=0)) + replicates = int(os.getenv('REPLICATES', default=REPLICATES)) + trigger_interval = float( + os.getenv( + 'TRIGGER_INTERVAL_SECONDS', + default=TRIGGER_INTERVAL_SECONDS + ) + ) + data_flush_interval = float( + os.getenv( + 'DATA_QUEUE_FLUSH_INTERVAL_SECONDS', + default=DATA_QUEUE_FLUSH_INTERVAL_SECONDS + ) + ) + data_queue: Queue = Queue() + # start main loop trigger_time = time.time() + data_flush_time = trigger_time + data_flush_interval while True: + # measure event if time.time() >= trigger_time: logging.debug('triggered event') # set new trigger time - trigger_interval = float( - os.getenv( - 'TRIGGER_INTERVAL_SECONDS', - default=TRIGGER_INTERVAL_SECONDS - ) - ) trigger_time += trigger_interval # run event - try: - event() - except Exception as e: - logging.error(e) - else: - logging.info('finished event') + event( + data_queue=data_queue, + user_id=user_id, + replicates=replicates + ) + # extra data flush + if ( + time.time() >= data_flush_time and + data_queue.qsize() > 0 + ): + logging.debug('triggered extra data flush') + # set new trigger time + data_flush_time += data_flush_interval + # flush data + flush_data( + data_queue=data_queue, + user_id=user_id + ) # polite pause time.sleep(1)