133 lines
3.5 KiB
Python
133 lines
3.5 KiB
Python
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 cron import cronjob
|
|
from bandwidth import measure as measure_bandwidth
|
|
from database import connect
|
|
|
|
|
|
# load default values
|
|
config = ConfigParser()
|
|
config.read(Path(__file__).parent / 'defaults.ini')
|
|
USER_ID = int(config.get('main', 'user_id'))
|
|
REPLICATES = int(config.get('main', 'replicates'))
|
|
SCHEDULE = config.get('cron', 'schedule')
|
|
TZ = config.get('cron', 'tz')
|
|
|
|
|
|
def flush_data(
|
|
data_queue: Queue,
|
|
user_id: int
|
|
) -> None:
|
|
# connect to database
|
|
db = connect()
|
|
# 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()
|
|
# 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:
|
|
data_queue.put(data) # put data back in queue
|
|
logging.error('failed sending data to database')
|
|
logging.debug(
|
|
'failed pushing to database:\n'
|
|
f'{json.dumps(payload_dict, indent=4)}\n'
|
|
f'with error: {e}'
|
|
)
|
|
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
|
|
flush_data(
|
|
data_queue=data_queue,
|
|
user_id=user_id
|
|
)
|
|
logging.debug('finished')
|
|
|
|
|
|
if __name__ == '__main__':
|
|
# setup
|
|
initialise_app()
|
|
data_queue: Queue = Queue()
|
|
user_id = int(os.getenv('USER_ID', default=USER_ID))
|
|
replicates = int(os.getenv('REPLICATES', default=REPLICATES))
|
|
schedule = os.getenv(
|
|
key='SCHEDULE',
|
|
default=SCHEDULE
|
|
)
|
|
tz = os.getenv(
|
|
key='TZ',
|
|
default=TZ
|
|
)
|
|
|
|
@cronjob(schedule=schedule, tz=tz)
|
|
def main_loop(
|
|
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
|
|
flush_data(
|
|
data_queue=data_queue,
|
|
user_id=user_id
|
|
)
|
|
logging.debug('finished')
|
|
# start main loop
|
|
main_loop(
|
|
data_queue=data_queue,
|
|
user_id=user_id,
|
|
replicates=replicates
|
|
)
|