+50
-45
@@ -6,6 +6,7 @@ from configparser import ConfigParser
|
||||
from pathlib import Path
|
||||
from queue import Queue
|
||||
from initialise_app import initialise_app
|
||||
from cron import cronjob
|
||||
from bandwidth import measure as measure_bandwidth
|
||||
from database import connect
|
||||
|
||||
@@ -14,18 +15,14 @@ from database import connect
|
||||
config = ConfigParser()
|
||||
config.read(Path(__file__).parent / 'defaults.ini')
|
||||
REPLICATES = int(config.get('main', 'replicates'))
|
||||
TRIGGER_INTERVAL_SECONDS = float(
|
||||
QUEUE_FLUSH_INTERVAL_SECONDS = float(
|
||||
config.get(
|
||||
'main',
|
||||
'trigger_interval_seconds'
|
||||
)
|
||||
)
|
||||
DATA_QUEUE_FLUSH_INTERVAL_SECONDS = float(
|
||||
config.get(
|
||||
'main',
|
||||
'data_queue_flush_interval_seconds'
|
||||
'queue_flush_interval_seconds'
|
||||
)
|
||||
)
|
||||
SCHEDULE = config.get('cron', 'schedule')
|
||||
TZ = config.get('cron', 'tz')
|
||||
|
||||
|
||||
def flush_data(
|
||||
@@ -34,8 +31,13 @@ def flush_data(
|
||||
) -> None:
|
||||
# connect to database
|
||||
db = connect()
|
||||
# flush queue
|
||||
# stop early
|
||||
queue_size = data_queue.qsize()
|
||||
if queue_size == 0:
|
||||
logging.debug('queue already empty')
|
||||
logging.debug('finished')
|
||||
return
|
||||
# flush queue
|
||||
for i in range(queue_size):
|
||||
# get data
|
||||
data = data_queue.get()
|
||||
@@ -48,13 +50,13 @@ def flush_data(
|
||||
try:
|
||||
db_id = db.insert_one(payload_dict).inserted_id
|
||||
except Exception as e:
|
||||
data_queue.put(data) # put data back in queue
|
||||
logging.error('failed sending data to database')
|
||||
logging.info(
|
||||
logging.debug(
|
||||
'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}')
|
||||
@@ -80,7 +82,6 @@ def event(
|
||||
# polite pause
|
||||
time.sleep(1)
|
||||
# flush queue
|
||||
if data_queue.qsize() > 0:
|
||||
flush_data(
|
||||
data_queue=data_queue,
|
||||
user_id=user_id
|
||||
@@ -91,48 +92,52 @@ def event(
|
||||
if __name__ == '__main__':
|
||||
# setup
|
||||
initialise_app()
|
||||
data_queue: Queue = Queue()
|
||||
user_id = int(os.getenv('USER_ID', default=0))
|
||||
replicates = int(os.getenv('REPLICATES', default=REPLICATES))
|
||||
trigger_interval = float(
|
||||
queue_flush_interval = float(
|
||||
os.getenv(
|
||||
'TRIGGER_INTERVAL_SECONDS',
|
||||
default=TRIGGER_INTERVAL_SECONDS
|
||||
'QUEUE_FLUSH_INTERVAL_SECONDS',
|
||||
default=QUEUE_FLUSH_INTERVAL_SECONDS
|
||||
)
|
||||
)
|
||||
data_flush_interval = float(
|
||||
os.getenv(
|
||||
'DATA_QUEUE_FLUSH_INTERVAL_SECONDS',
|
||||
default=DATA_QUEUE_FLUSH_INTERVAL_SECONDS
|
||||
schedule = os.getenv(
|
||||
key='SCHEDULE',
|
||||
default=SCHEDULE
|
||||
)
|
||||
tz = os.getenv(
|
||||
key='TZ',
|
||||
default=TZ
|
||||
)
|
||||
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
|
||||
|
||||
@cronjob(schedule=schedule, tz=tz)
|
||||
def main_loop(
|
||||
data_queue: Queue,
|
||||
user_id: int,
|
||||
replicates: int = REPLICATES
|
||||
) -> None:
|
||||
# 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
|
||||
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
|
||||
flush_data(
|
||||
data_queue=data_queue,
|
||||
user_id=user_id
|
||||
)
|
||||
# polite pause
|
||||
time.sleep(1)
|
||||
logging.debug('finished')
|
||||
# start main loop
|
||||
main_loop(
|
||||
data_queue=data_queue,
|
||||
user_id=user_id,
|
||||
replicates=replicates
|
||||
)
|
||||
|
||||
Reference in New Issue
Block a user