Files
bandwidth_probing/code/main.py
T
2023-07-24 10:13:49 +02:00

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
)