From fd50bbbab7cccbdd71aa8ee77be3c194323d7f42 Mon Sep 17 00:00:00 2001 From: Bastian Wagner Date: Sat, 15 Aug 2026 16:14:25 +0200 Subject: [PATCH] feat: schedule periodic user synchronization --- app/main.py | 38 +++++++++++++++++++++++++++++++- app/sync/scheduler.py | 42 ++++++++++++++++++++++++++++++++++++ tests/sync/test_scheduler.py | 42 ++++++++++++++++++++++++++++++++++++ tests/test_main_lifespan.py | 30 ++++++++++++++++++++++++++ 4 files changed, 151 insertions(+), 1 deletion(-) create mode 100644 app/sync/scheduler.py create mode 100644 tests/sync/test_scheduler.py create mode 100644 tests/test_main_lifespan.py diff --git a/app/main.py b/app/main.py index 6827988..4223c19 100644 --- a/app/main.py +++ b/app/main.py @@ -1,8 +1,16 @@ +from contextlib import asynccontextmanager + from fastapi import FastAPI from starlette.middleware.sessions import SessionMiddleware from app.config import Settings, get_settings from app.db.session import create_db_engine, create_session_factory, initialize_schema +from app.fit.rewriter import convert_fit_device +from app.garmin.uploader import GarminUploader +from app.mywhoosh.client import MyWhooshClient +from app.security.credentials import CredentialCipher +from app.sync.manager import SyncManager +from app.sync.scheduler import SyncScheduler from app.web.routes import router as web_router @@ -12,7 +20,35 @@ def create_app(settings: Settings | None = None) -> FastAPI: resolved.tokens_dir.mkdir(parents=True, exist_ok=True) resolved.activities_dir.mkdir(parents=True, exist_ok=True) - app = FastAPI(title="MyWhoosh Garmin Sync") + @asynccontextmanager + async def lifespan(app: FastAPI): + cipher = CredentialCipher(resolved.credential_encryption_key) + + def mywhoosh_factory(token_store): + return MyWhooshClient(token_store) + + def garmin_factory(email, password, tokenstore): + return GarminUploader(email=email, password=password, tokenstore=tokenstore) + + sync_manager = SyncManager( + session_factory=app.state.session_factory, + credential_cipher=cipher, + settings=resolved, + mywhoosh_factory=mywhoosh_factory, + garmin_factory=garmin_factory, + fit_converter=convert_fit_device, + ) + app.state.sync_manager = sync_manager + + scheduler = SyncScheduler(sync_manager, interval_seconds=resolved.sync_interval_minutes * 60) + app.state.scheduler = scheduler + await scheduler.start() + + yield + + await scheduler.stop() + + app = FastAPI(title="MyWhoosh Garmin Sync", lifespan=lifespan) app.state.settings = resolved engine = create_db_engine(resolved.database_url) diff --git a/app/sync/scheduler.py b/app/sync/scheduler.py new file mode 100644 index 0000000..3194a9b --- /dev/null +++ b/app/sync/scheduler.py @@ -0,0 +1,42 @@ +import asyncio +import logging +from datetime import datetime, timedelta, timezone + +logger = logging.getLogger(__name__) + + +class SyncScheduler: + def __init__(self, manager, *, interval_seconds: float) -> None: + self.manager = manager + self.interval_seconds = interval_seconds + self._task: asyncio.Task | None = None + self._stop = asyncio.Event() + self.last_tick = None + self.next_tick = None + + async def run_once(self) -> None: + self.last_tick = datetime.now(timezone.utc) + try: + await self.manager.sync_all_enabled() + except Exception: + logger.exception("sync_all_enabled failed during scheduled tick") + finally: + self.next_tick = datetime.now(timezone.utc) + timedelta(seconds=self.interval_seconds) + + async def _run(self) -> None: + while not self._stop.is_set(): + await self.run_once() + try: + await asyncio.wait_for(self._stop.wait(), timeout=self.interval_seconds) + except TimeoutError: + pass + + async def start(self) -> None: + self._stop.clear() + self._task = asyncio.create_task(self._run()) + + async def stop(self) -> None: + self._stop.set() + if self._task is not None: + await self._task + self._task = None diff --git a/tests/sync/test_scheduler.py b/tests/sync/test_scheduler.py new file mode 100644 index 0000000..9843ba9 --- /dev/null +++ b/tests/sync/test_scheduler.py @@ -0,0 +1,42 @@ +import asyncio + +import pytest + +from app.sync.scheduler import SyncScheduler + + +class FakeSyncManager: + def __init__(self, results=None) -> None: + self.results = results if results is not None else [] + self.calls = 0 + + async def sync_all_enabled(self): + self.calls += 1 + if self.results and isinstance(self.results[0], Exception): + raise self.results[0] + return list(self.results) + + +@pytest.mark.asyncio +async def test_scheduler_calls_sync_all_and_survives_failure() -> None: + fake = FakeSyncManager(results=[RuntimeError("one user failed")]) + scheduler = SyncScheduler(fake, interval_seconds=0.01) + await scheduler.start() + await asyncio.sleep(0.035) + await scheduler.stop() + assert fake.calls >= 2 + assert scheduler.last_tick is not None + + +@pytest.mark.asyncio +async def test_scheduler_stop_cancels_the_loop() -> None: + fake = FakeSyncManager() + scheduler = SyncScheduler(fake, interval_seconds=0.01) + await scheduler.start() + await asyncio.sleep(0.035) + await scheduler.stop() + + assert scheduler._task is None + calls_after_stop = fake.calls + await asyncio.sleep(0.05) + assert fake.calls == calls_after_stop diff --git a/tests/test_main_lifespan.py b/tests/test_main_lifespan.py new file mode 100644 index 0000000..d2dbca3 --- /dev/null +++ b/tests/test_main_lifespan.py @@ -0,0 +1,30 @@ +from pathlib import Path + +from cryptography.fernet import Fernet +from fastapi.testclient import TestClient + +from app.config import Settings +from app.main import create_app +from app.sync.manager import SyncManager + + +def test_lifespan_wires_sync_manager_and_scheduler(tmp_path: Path) -> None: + settings = Settings( + ADMIN_PASSWORD="admin-secret", + SECRET_KEY="0123456789abcdef0123456789abcdef", + CREDENTIAL_ENCRYPTION_KEY=Fernet.generate_key().decode("ascii"), + DATA_DIR=str(tmp_path), + DATABASE_URL=f"sqlite:///{tmp_path / 'app.db'}", + SYNC_INTERVAL_MINUTES=5, + ) + app = create_app(settings) + + # No users are seeded in this database, so sync_all_enabled() has nothing + # to iterate over and the real MyWhoosh/Garmin factories are never invoked. + with TestClient(app): + assert app.state.sync_manager is not None + assert isinstance(app.state.sync_manager, SyncManager) + assert app.state.scheduler is not None + assert app.state.scheduler.last_tick is not None + + app.state.db_engine.dispose()