Isolated the literal Sync API function to its own module, and streamlined the sync UI feedback
This commit is contained in:
parent
dc46d4136a
commit
28d4f3d276
5 changed files with 158 additions and 163 deletions
|
|
@ -1,5 +1,6 @@
|
|||
# new sync refactor:
|
||||
from core.services.sync.sync_service import coordinate_cache_sync
|
||||
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.errors.logger import logger
|
||||
# prior versions:
|
||||
from core.Constants import Constants
|
||||
|
|
@ -48,16 +49,16 @@ class ClientController:
|
|||
@staticmethod
|
||||
def sync(client_observer: ClientObserver = None, connection_observer: ConnectionObserver = None):
|
||||
|
||||
from core.models.manage.session_management import init_session, close_session
|
||||
if client_observer is not None:
|
||||
client_observer.notify('synchronizing', "Fetching list of new data ..")
|
||||
|
||||
# Use Cached Method:
|
||||
init_session()
|
||||
|
||||
result = coordinate_cache_sync(ClientObserver, ConnectionObserver)
|
||||
|
||||
logger.info(f"We got a Result from the API of {result}")
|
||||
|
||||
close_session()
|
||||
|
||||
# Outright Error:
|
||||
if not result["success"]:
|
||||
error_msg = result["error"]
|
||||
|
|
@ -74,14 +75,42 @@ class ClientController:
|
|||
|
||||
# We only make it past this point if there's New Data
|
||||
|
||||
# Fetch and update those models...
|
||||
# flag for after the save,
|
||||
data_was_saved = False
|
||||
|
||||
try:
|
||||
# Fetch and update the real data (no longer metadata)...
|
||||
from core.controllers.ConnectionController import ConnectionController
|
||||
ConnectionController.with_preferred_connection(task=ClientController.__sync, changed_tables=changed_tables, client_observer=client_observer, connection_observer=connection_observer)
|
||||
|
||||
# We set the flag to true,
|
||||
# the reason we use a flag, and don't just save it right here,
|
||||
# is because we want to isolate the success (or failure) of the real data,
|
||||
# from the potential failure of the ORM session metadata.
|
||||
data_was_saved = True
|
||||
|
||||
except:
|
||||
# sync failed here,
|
||||
if client_observer is not None:
|
||||
client_observer.notify('synchronizing', 'Error! Sync Failed.')
|
||||
client_observer.notify('synchronizing', 'Fetch Failed, but you can use old data.')
|
||||
finally:
|
||||
if data_was_saved:
|
||||
filtered_metadata = result["filtered_metadata"] # from the top of the function
|
||||
save_successful = save_metadata(filtered_metadata) # the "save_data" function is inside sync_service
|
||||
|
||||
# Regardless of the outcome,
|
||||
close_session()
|
||||
|
||||
if client_observer is None:
|
||||
logger.error("Error: No client_observer to update the UI, the final part of the sync function skipped")
|
||||
return # can't update their UI
|
||||
|
||||
if save_successful:
|
||||
logger.info("Metadata Saved Successfully")
|
||||
client_observer.notify('synchronized', "Fetch & Save Complete!")
|
||||
else:
|
||||
client_observer.notify('synchronizing', "Saving List of Metadata Failed.")
|
||||
|
||||
|
||||
@staticmethod
|
||||
def update(client_observer: ClientObserver = None, connection_observer: ConnectionObserver = None):
|
||||
|
|
@ -101,56 +130,51 @@ class ClientController:
|
|||
@staticmethod
|
||||
def __sync(changed_tables: list, client_observer: Optional[ClientObserver] = None, proxies: Optional[dict] = None):
|
||||
|
||||
if client_observer is not None:
|
||||
client_observer.notify('synchronizing')
|
||||
|
||||
if "applications" in changed_tables:
|
||||
logger.info("Sync applications..")
|
||||
if client_observer is not None:
|
||||
client_observer.notify('synchronizing', 'Fetching Browser List')
|
||||
client_observer.notify('synchronizing', 'Fetching Browser List..')
|
||||
# noinspection PyProtectedMember
|
||||
ApplicationController._sync(proxies=proxies)
|
||||
|
||||
if "application_versions" in changed_tables:
|
||||
logger.info("Sync Application Versions..")
|
||||
if client_observer is not None:
|
||||
client_observer.notify('synchronizing', 'Fetching Browser Version List')
|
||||
client_observer.notify('synchronizing', 'Fetching Browser Version List..')
|
||||
# noinspection PyProtectedMember
|
||||
ApplicationVersionController._sync(proxies=proxies)
|
||||
|
||||
if "client_version" in changed_tables:
|
||||
logger.info("Sync of client version")
|
||||
if client_observer is not None:
|
||||
client_observer.notify('synchronizing', 'Fetching Client Version List')
|
||||
client_observer.notify('synchronizing', 'Fetching Client Version List..')
|
||||
# noinspection PyProtectedMember
|
||||
ClientVersionController._sync(proxies=proxies)
|
||||
|
||||
if "operators" in changed_tables:
|
||||
logger.info("Sync of Operators")
|
||||
if client_observer is not None:
|
||||
client_observer.notify('synchronizing', 'Fetching Operators List')
|
||||
client_observer.notify('synchronizing', 'Fetching Operators List..')
|
||||
# noinspection PyProtectedMember
|
||||
OperatorController._sync(proxies=proxies)
|
||||
|
||||
if "locations" in changed_tables:
|
||||
logger.info("Sync of Locations")
|
||||
if client_observer is not None:
|
||||
client_observer.notify('synchronizing', 'Fetching Locations List')
|
||||
client_observer.notify('synchronizing', 'Fetching Locations List..')
|
||||
# noinspection PyProtectedMember
|
||||
LocationController._sync(proxies=proxies)
|
||||
|
||||
if "subscriptions" in changed_tables:
|
||||
logger.info("Sync of Subscriptions")
|
||||
if client_observer is not None:
|
||||
client_observer.notify('synchronizing', 'Fetching Subscription List')
|
||||
client_observer.notify('synchronizing', 'Fetching Subscription List..')
|
||||
# noinspection PyProtectedMember
|
||||
SubscriptionPlanController._sync(proxies=proxies)
|
||||
|
||||
ConfigurationController.update_last_synced_at()
|
||||
|
||||
logger.info("Entire Sync Completed Successfully")
|
||||
if client_observer is not None:
|
||||
client_observer.notify('synchronized')
|
||||
logger.info("Real Data Fetch Completed Successfully")
|
||||
|
||||
@staticmethod
|
||||
def __update(client_observer: Optional[ClientObserver] = None, proxies: Optional[dict] = None):
|
||||
|
|
|
|||
77
core/services/sync/get_metadata_from_api.py
Normal file
77
core/services/sync/get_metadata_from_api.py
Normal file
|
|
@ -0,0 +1,77 @@
|
|||
from core.Constants import Constants
|
||||
from core.observers.BaseObserver import BaseObserver
|
||||
from core.services.networking.get_data_from_server import get_data_from_server
|
||||
from core.errors.logger import logger
|
||||
from core.observers.ClientObserver import ClientObserver
|
||||
from core.observers.ConnectionObserver import ConnectionObserver
|
||||
from typing import Optional
|
||||
|
||||
|
||||
# this can be spoofed with the isolation testing comment below
|
||||
|
||||
def get_metadata_from_api(
|
||||
client_observer: Optional[ClientObserver] = None,
|
||||
connection_observer: Optional[ConnectionObserver] = None
|
||||
) -> dict:
|
||||
logger.info("Syncing with the API...")
|
||||
|
||||
# Get and validate the base URL
|
||||
rejected_list = [None, False, ""]
|
||||
base_url = Constants.SP_API_BASE_URL
|
||||
if base_url in rejected_list:
|
||||
error_msg = "Invalid base URL from configuration"
|
||||
logger.error(error_msg)
|
||||
return {"success": False, "data": None, "error": error_msg}
|
||||
|
||||
# Construct the full endpoint URL
|
||||
url = f"{base_url}/cachedsync"
|
||||
logger.debug(f"API endpoint: {url}")
|
||||
|
||||
try:
|
||||
logger.debug("Fetching data from API server...")
|
||||
sync_results = get_data_from_server(url, connection_observer)
|
||||
|
||||
if sync_results in rejected_list:
|
||||
error_msg = "API returned no data"
|
||||
logger.error(error_msg)
|
||||
return {"success": False, "data": None, "error": error_msg}
|
||||
|
||||
logger.debug(f"Raw API response: {sync_results}")
|
||||
|
||||
final_result = sync_results.get("data", None)
|
||||
|
||||
if final_result:
|
||||
logger.info("Successfully retrieved sync versions metadata from API")
|
||||
return {"success": True, "data": final_result, "error": None}
|
||||
|
||||
error_msg = "API response missing 'data' field"
|
||||
logger.error(error_msg)
|
||||
return {"success": False, "data": None, "error": error_msg}
|
||||
|
||||
except Exception as e:
|
||||
error_msg = f"API fetch failed: {str(e)}"
|
||||
logger.error(error_msg)
|
||||
return {"success": False, "data": None, "error": error_msg}
|
||||
|
||||
|
||||
|
||||
|
||||
|
||||
|
||||
# ============== ISOLATION TESTING ==================
|
||||
# SPOOF ISOLATION TEST:
|
||||
# def _get_sync_cache_from_api(
|
||||
# client_observer: Optional[ClientObserver] = None,
|
||||
# connection_observer: Optional[ConnectionObserver] = None
|
||||
# ) -> dict:
|
||||
# final_result = {
|
||||
# "version": 9,
|
||||
# "applications": 5,
|
||||
# "application_versions": 4,
|
||||
# "client_version": 3,
|
||||
# "operators": 2,
|
||||
# "locations": 3,
|
||||
# "subscriptions": 1
|
||||
# }
|
||||
# return {"success": True, "data": final_result, "error": None}
|
||||
|
||||
|
|
@ -1,33 +1,27 @@
|
|||
# from __future__ import annotations
|
||||
# from typing import TYPE_CHECKING
|
||||
|
||||
# if TYPE_CHECKING:
|
||||
# from core.essentials.observers.ConnectionObserver import ConnectionObserver
|
||||
# from core.observers.TicketObserver import TicketObserver
|
||||
|
||||
# comparisons:
|
||||
# comparisons & api calls:
|
||||
from core.services.sync.compare_tables import compare_tables
|
||||
from core.services.sync.get_metadata_from_api import get_metadata_from_api
|
||||
|
||||
# unique database:
|
||||
# ORM for metadata:
|
||||
from core.models.manage.session_management import get_session
|
||||
from core.models.manage.get_from_model import get_from_model
|
||||
from core.models.manage.insert import insert_into_model
|
||||
from core.models.orm_models.CachedSync import CachedSync
|
||||
|
||||
# usual core infrastructure:
|
||||
from core.Constants import Constants
|
||||
from core.observers.BaseObserver import BaseObserver
|
||||
from core.services.networking.get_data_from_server import get_data_from_server
|
||||
from core.errors.logger import logger
|
||||
from core.observers.BaseObserver import BaseObserver
|
||||
from core.observers.ClientObserver import ClientObserver
|
||||
from core.observers.ConnectionObserver import ConnectionObserver
|
||||
|
||||
# generic
|
||||
from sqlalchemy import func
|
||||
from typing import Optional
|
||||
from typing import Optional
|
||||
|
||||
|
||||
def _get_cached_metadata():
|
||||
def _get_cached_metadata() -> tuple:
|
||||
"""Retrieve cached metadata from the database.
|
||||
|
||||
Returns:
|
||||
|
|
@ -54,7 +48,7 @@ def _get_cached_metadata():
|
|||
return old_data, previous_highest
|
||||
|
||||
|
||||
def _get_changed_models(new_data, old_data):
|
||||
def _get_changed_models(new_data: dict, old_data: dict) -> tuple:
|
||||
"""Determine which models changed between API and cached versions.
|
||||
|
||||
Args:
|
||||
|
|
@ -80,7 +74,7 @@ def _get_changed_models(new_data, old_data):
|
|||
return changed_tables, new_model_types
|
||||
|
||||
|
||||
def _filter_valid_models(data, invalid_keys):
|
||||
def _filter_valid_models(data: dict, invalid_keys: list) -> dict:
|
||||
"""Remove invalid model types from data dict.
|
||||
|
||||
Args:
|
||||
|
|
@ -97,11 +91,19 @@ def _filter_valid_models(data, invalid_keys):
|
|||
return filtered
|
||||
|
||||
|
||||
def save_data(data: dict) -> bool:
|
||||
"""Save metadata to database."""
|
||||
def save_metadata(metadata: dict) -> bool:
|
||||
"""
|
||||
Save metadata to the database using the ORM.
|
||||
|
||||
Called By:
|
||||
ClientController.sync
|
||||
|
||||
Returns:
|
||||
Boolean of if it worked or not.
|
||||
"""
|
||||
try:
|
||||
logger.debug(f"Inserting metadata into database: {data}")
|
||||
did_it_work = insert_into_model(CachedSync, data, True)
|
||||
logger.debug(f"Inserting metadata into database: {metadata}")
|
||||
did_it_work = insert_into_model(CachedSync, metadata, True)
|
||||
|
||||
if did_it_work:
|
||||
logger.debug("Metadata saved successfully")
|
||||
|
|
@ -115,7 +117,7 @@ def save_data(data: dict) -> bool:
|
|||
return False
|
||||
|
||||
|
||||
def _log_what_changed(new_model_types, changed_tables):
|
||||
def _log_what_changed(new_model_types: list, changed_tables: list):
|
||||
if new_model_types:
|
||||
logger.debug(f"New model types detected (reserved): {new_model_types}")
|
||||
|
||||
|
|
@ -133,17 +135,21 @@ def coordinate_cache_sync(
|
|||
"""
|
||||
Fetch version metadata from API and return changed/new model keys.
|
||||
|
||||
Called by:
|
||||
ClientController.sync
|
||||
|
||||
Returns:
|
||||
Dictionary with keys:
|
||||
- 'success' (bool): Whether comparison succeeded
|
||||
- 'changed_tables' (list): Model keys that changed
|
||||
- 'new_model_types' (list): New model keys (reserved for future use)
|
||||
- 'error' (str or None): Error message if success is False
|
||||
- 'filtered_metadata': the metadata from the server after removing keys not in the ORM
|
||||
"""
|
||||
|
||||
api_result = _get_sync_cache_from_api(ClientObserver, ConnectionObserver)
|
||||
api_result = get_metadata_from_api(ClientObserver, ConnectionObserver)
|
||||
|
||||
# we trust the '_get_sync_cache_from_api' function to give us a dictionary:
|
||||
# we trust the 'get_metadata_from_api' function to give us a dictionary:
|
||||
if not api_result.get("success"):
|
||||
error_msg = api_result.get("error", "Unknown API error")
|
||||
logger.error(f"API sync failed: {error_msg}")
|
||||
|
|
@ -223,89 +229,17 @@ def coordinate_cache_sync(
|
|||
# this is not saving it to SQL, it's a temp debug log:
|
||||
_log_what_changed(new_model_types, changed_tables)
|
||||
|
||||
# Filter out new model types for the upcoming save:
|
||||
# Filter out new model types for the upcoming save LATER:
|
||||
filtered_data = _filter_valid_models(new_data, new_model_types)
|
||||
|
||||
# Save only currently existing models to SQL:
|
||||
save_successful = save_data(filtered_data)
|
||||
if save_successful:
|
||||
return {
|
||||
"success": True,
|
||||
"changed_tables": changed_tables,
|
||||
"new_model_types": new_model_types,
|
||||
"error": None,
|
||||
}
|
||||
else:
|
||||
return {
|
||||
"success": False,
|
||||
"changed_tables": [],
|
||||
"new_model_types": [],
|
||||
"error": "Failed to save metadata to database",
|
||||
}
|
||||
# we'll save the filtered metadata in the main module after the real sync is successful,
|
||||
# because if we don't (and save it right now), then if the upcoming sync fails,
|
||||
# we'd have wrong metadata saved (without matching real data) and skip future attempts.
|
||||
return {
|
||||
"success": True,
|
||||
"changed_tables": changed_tables,
|
||||
"new_model_types": new_model_types,
|
||||
"error": None,
|
||||
"filtered_metadata": filtered_data
|
||||
}
|
||||
|
||||
|
||||
|
||||
|
||||
def _get_sync_cache_from_api(
|
||||
client_observer: Optional[ClientObserver] = None,
|
||||
connection_observer: Optional[ConnectionObserver] = None
|
||||
) -> dict:
|
||||
logger.info("Syncing with the API...")
|
||||
|
||||
# Get and validate the base URL
|
||||
rejected_list = [None, False, ""]
|
||||
base_url = Constants.SP_API_BASE_URL
|
||||
if base_url in rejected_list:
|
||||
error_msg = "Invalid base URL from configuration"
|
||||
logger.error(error_msg)
|
||||
return {"success": False, "data": None, "error": error_msg}
|
||||
|
||||
# Construct the full endpoint URL
|
||||
url = f"{base_url}/cachedsync"
|
||||
logger.debug(f"API endpoint: {url}")
|
||||
|
||||
try:
|
||||
logger.debug("Fetching data from API server...")
|
||||
sync_results = get_data_from_server(url, connection_observer)
|
||||
|
||||
if sync_results in rejected_list:
|
||||
error_msg = "API returned no data"
|
||||
logger.error(error_msg)
|
||||
return {"success": False, "data": None, "error": error_msg}
|
||||
|
||||
logger.debug(f"Raw API response: {sync_results}")
|
||||
|
||||
final_result = sync_results.get("data", None)
|
||||
|
||||
if final_result:
|
||||
logger.info("Successfully retrieved sync versions metadata from API")
|
||||
return {"success": True, "data": final_result, "error": None}
|
||||
|
||||
error_msg = "API response missing 'data' field"
|
||||
logger.error(error_msg)
|
||||
return {"success": False, "data": None, "error": error_msg}
|
||||
|
||||
except Exception as e:
|
||||
error_msg = f"API fetch failed: {str(e)}"
|
||||
logger.error(error_msg)
|
||||
return {"success": False, "data": None, "error": error_msg}
|
||||
|
||||
|
||||
|
||||
|
||||
|
||||
|
||||
# ============== ISOLATION TESTING ==================
|
||||
|
||||
# SPOOF ISOLATION TEST:
|
||||
# def _get_sync_cache_from_api(
|
||||
# final_result = {
|
||||
# "version": 5,
|
||||
# "applications": 3,
|
||||
# "application_versions": 2,
|
||||
# "client_version": 2,
|
||||
# "operators": 2,
|
||||
# "locations": 2,
|
||||
# "subscriptions": 1
|
||||
# }
|
||||
# return {"success": True, "data": final_result, "error": None}
|
||||
|
|
|
|||
|
|
@ -1,40 +0,0 @@
|
|||
from core.controllers.ClientController import ClientController
|
||||
from core.observers.ClientObserver import ClientObserver
|
||||
from core.observers.ConnectionObserver import ConnectionObserver
|
||||
from core.models.manage.session_management import init_session, close_session
|
||||
|
||||
class PrototypeSync():
|
||||
def __init__(self, parent=None):
|
||||
client_observer = ClientObserver()
|
||||
connection_observer = ConnectionObserver()
|
||||
client_observer.subscribe('synchronizing', lambda event: print('Synchronizing...\n'))
|
||||
client_observer.subscribe('updating', lambda event: print('Updating client...'))
|
||||
client_observer.subscribe('update_progressing', lambda event: print(f'Current progress: {event.meta.get('progress'):.2f}%', flush=True, end='\r'))
|
||||
client_observer.subscribe('updated', lambda event: print('\n'))
|
||||
|
||||
connection_observer.subscribe('connecting', lambda event: self.update_status(
|
||||
f'[{event.subject.get("attempt_count")}/{event.subject.get("maximum_number_of_attempts")}] Performing connection attempt...'))
|
||||
|
||||
connection_observer.subscribe('tor_bootstrapping', lambda event: self.update_status(
|
||||
'Establishing Tor connection...'))
|
||||
|
||||
connection_observer.subscribe('tor_bootstrap_progressing', lambda event: self.update_status(
|
||||
f'Bootstrapping Tor {event.meta.get('progress'):.2f}%'))
|
||||
|
||||
connection_observer.subscribe(
|
||||
'tor_bootstrapped', lambda event: self.update_status('Tor connection established.'))
|
||||
|
||||
|
||||
init_session()
|
||||
|
||||
print("Sync..")
|
||||
ClientController.sync(client_observer=client_observer, connection_observer=connection_observer)
|
||||
print("Done")
|
||||
|
||||
# On app shutdown
|
||||
close_session()
|
||||
|
||||
|
||||
def update_status(self, text, clear=False):
|
||||
if text:
|
||||
print(text)
|
||||
|
|
@ -1,6 +1,6 @@
|
|||
[project]
|
||||
name = "sp-hydra-veil-core"
|
||||
version = "2.3.5"
|
||||
version = "2.3.6"
|
||||
authors = [
|
||||
{ name = "Simplified Privacy" },
|
||||
]
|
||||
|
|
|
|||
Loading…
Reference in a new issue