sp-hydra-veil-core/core/controllers/SyncController.py

217 lines
7.8 KiB
Python

from core.services.sync import legacy_insert
from core.models.manage.session_management import init_session, close_session
from core.services.sync.sync_service import coordinate_cache_sync, save_metadata
from core.services.sync.insert_for_orm import insert_one_orm_model
from core.services.networking.httpx import connect
# return objects
from core.services.networking.api_requests.ApiResponseModel import ApiResponse, ErrorType
from core.models.DatabaseOperation import DatabaseOperation, DBErrorType
from core.models.Result import Result, ResultError
from core.errors.logger import logger
# ORM models that can be sync'ed:
from core.models.orm_models.Location import Location
from core.models.orm_models.Operator import Operator
from core.models.orm_models.Dependency import Dependency
from core.Constants import Constants
from core.controllers.ApplicationController import ApplicationController
from core.controllers.ApplicationVersionController import ApplicationVersionController
from core.controllers.ClientVersionController import ClientVersionController
from core.controllers.ConfigurationController import ConfigurationController
from core.controllers.SubscriptionPlanController import SubscriptionPlanController
from core.observers.ClientObserver import ClientObserver
from core.observers.ConnectionObserver import ConnectionObserver
from core.models.orm_models.ApplicationVersion import ApplicationVersion
# generic
import sys
import json
ORM_TABLES = {
"locations": Location,
"operators": Operator,
"application_versions": ApplicationVersion,
"dependencies": Dependency
}
LEGACY_SQL_FUNCT_DICT = {
"applications": legacy_insert.for_applications,
"client_version": legacy_insert.for_client_version,
"subscriptions": legacy_insert.for_subscriptions,
}
APP_CODES = {
"firefox": 1,
"chromium": 2,
"brave": 3,
"librewolf": 5
}
def call_legacy_insert_function(key: str, new_data: dict):
func = LEGACY_SQL_FUNCT_DICT.get(key)
if func is None:
return DatabaseOperation(valid=False, error_type=DBErrorType.UNKNOWN_MODEL)
if key in APP_CODES:
code = APP_CODES[key]
return func(new_data, code)
else:
return func(new_data)
def get_orm_model(key) -> bool:
if key not in ORM_TABLES:
return False
return ORM_TABLES[key]
def new_sync(client_observer: ClientObserver, connection_observer: ConnectionObserver) -> Result:
client_observer.notify('synchronizing', "Fetching list of new data ..")
####################################
# METADATA. Should we even sync?
####################################
metadata_result = coordinate_cache_sync(client_observer, connection_observer)
# Outright Error:
if not metadata_result["success"]:
error_msg = metadata_result["error"]
client_observer.notify('synchronizing', f'Error! {error_msg}')
return Result(valid=False)
# Nothing changed if the 'changed_tables' variable does NOT exist
changed_tables = metadata_result["changed_tables"]
# NEW MODELS:
new_model_types = metadata_result.get("new_model_types", [])
if new_model_types:
for each_new_model in new_model_types:
if each_new_model in ORM_TABLES:
# Then we know how to handle this, let's sync it:
logger.info(f"Adding the brand new {each_new_model} to the changed tables.")
changed_tables.append(each_new_model)
else:
logger.info(f"Skipping {each_new_model} because we don't know how to handle it yet.")
# CHECK ON CHANGED TABLES
if not changed_tables:
client_observer.notify('synchronized')
return Result(valid=True)
# We only make it past this point if there's New Data
####################################
# API CALLS: GET NEW DATA IN BULK
####################################
results = connect.bulk_async(
wanted_list=changed_tables,
observer=connection_observer,
client_observer=client_observer
)
if not results.valid:
error_msg = f"Sync failed! {results.error_type}"
logger.error(f"{error_msg} {results.message}")
client_observer.notify('synchronizing', f'Error! {error_msg}')
return Result(valid=False)
quantity_of_entries = len(results.data)
logger.info(f"We have valid API call results. There are {quantity_of_entries} entries")
all_results = results.data
####################################
# LOOP INSERT INTO DATABASE
####################################
client_observer.notify('synchronizing', f'Inserting into Database..')
skipped = []
for key, each_api_result in all_results.items():
if not each_api_result.valid:
logger.info(f"Skipping invalid api response for {key}")
skipped.append(key)
continue
logger.info(f"Inserting valid api data for {key} into the Database")
each_insert = insert_data(
key=key,
each_api_calls_data=each_api_result.data,
client_observer=client_observer
)
logger.info(f"Exited the insert data function with a result of {each_insert.valid}")
if each_insert.valid:
logger.info(f"Success with insert of {key}")
client_observer.notify('synchronizing', f'Inserted {key}')
continue
else:
skipped.append(key)
logger.error(f"Database error with inserting {key}, because {each_insert.error_type}")
client_observer.notify('synchronizing', f'Failed inserting {key} because {each_insert.message}')
logger.info("This only gets triggered in errors, but Moving on to the next item..")
####################################
# FINAL EVALUATION
####################################
total_skipped = len(skipped)
logger.info(f"We exited the loop, doing the final evaluation. And a total of {total_skipped} were skipped.")
ConfigurationController.update_last_synced_at()
filtered_metadata = metadata_result["filtered_metadata"] # from the top of the function
if total_skipped == 0:
client_observer.notify('synchronized', "Fetch & Save Complete!")
save_successful = save_metadata(filtered_metadata)
if save_successful:
return Result(valid=True, message="Finshed sync.")
else:
error_msg = "Finshed sync, but had issues with the saving of metadata for next time."
logger.error(error_msg)
return Result(valid=True, message=error_msg)
elif total_skipped < quantity_of_entries:
error_msg = f"Partial Success. {total_skipped} skipped."
client_observer.notify('synchronized', error_msg)
return Result(valid=True, data=skipped, message=error_msg)
else:
error_msg = f"Sync Failed. All {total_skipped} entries were skipped!"
client_observer.notify('synchronized', error_msg)
return Result(valid=False, data=skipped, message="Complete Failure, all data failed to insert.")
def insert_data(
key: str,
each_api_calls_data: dict,
client_observer: ClientObserver,
) -> DatabaseOperation:
####################################
# NEW ORM SYSTEM
####################################
new_orm_model = get_orm_model(key)
if new_orm_model:
db_result = insert_one_orm_model(
which_key=key,
which_model=new_orm_model,
new_data=each_api_calls_data,
override=True
)
return db_result
####################################
# LEGACY MANUAL SQL
####################################
return call_legacy_insert_function(key=key, new_data=each_api_calls_data)
####################################
# UNKNOWN MODEL
####################################
return DatabaseOperation(valid=False, error_type=DBErrorType.UNKNOWN_MODEL)