import time import logging import os import json from configparser import ConfigParser from pathlib import Path from queue import Queue from initialise_app import initialise_app from bandwidth import measure as measure_bandwidth from database import connect # load default values config = ConfigParser() config.read(Path(__file__).parent / 'defaults.ini') REPLICATES = int(config.get('main', 'replicates')) 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 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() # prepare payload payload_dict = { 'user_id': user_id, 'data': data } # send to database try: db_id = db.insert_one(payload_dict).inserted_id except Exception as e: logging.error('failed sending data to database') logging.info( 'failed pushing to database:\n' f'{json.dumps(payload_dict, indent=4)}\n' f'with error: {e}' ) data_queue.put(data) # put data back in queue continue else: logging.debug(f'data sent to database received id: {db_id}') 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_time += trigger_interval # run 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)