feat: add SQLite persistence models
This commit is contained in:
0
app/db/__init__.py
Normal file
0
app/db/__init__.py
Normal file
92
app/db/models.py
Normal file
92
app/db/models.py
Normal file
@@ -0,0 +1,92 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import enum
|
||||
from datetime import datetime, timezone
|
||||
|
||||
from sqlalchemy import Boolean, DateTime, Enum, ForeignKey, Integer, String, Text, UniqueConstraint
|
||||
from sqlalchemy.orm import DeclarativeBase, Mapped, mapped_column, relationship
|
||||
|
||||
|
||||
def utcnow() -> datetime:
|
||||
return datetime.now(timezone.utc)
|
||||
|
||||
|
||||
class Base(DeclarativeBase):
|
||||
pass
|
||||
|
||||
|
||||
class HealthState(str, enum.Enum):
|
||||
HEALTHY = "healthy"
|
||||
SYNCING = "syncing"
|
||||
DEGRADED = "degraded"
|
||||
ACTION_REQUIRED = "action_required"
|
||||
DISABLED = "disabled"
|
||||
|
||||
|
||||
class ActivityStatus(str, enum.Enum):
|
||||
DISCOVERED = "discovered"
|
||||
DOWNLOADED = "downloaded"
|
||||
CONVERTED = "converted"
|
||||
IMPORTED = "imported"
|
||||
DUPLICATE = "duplicate"
|
||||
FAILED = "failed"
|
||||
|
||||
|
||||
class SyncRunStatus(str, enum.Enum):
|
||||
RUNNING = "running"
|
||||
SUCCESS = "success"
|
||||
PARTIAL = "partial"
|
||||
FAILED = "failed"
|
||||
|
||||
|
||||
class SyncUser(Base):
|
||||
__tablename__ = "sync_users"
|
||||
|
||||
id: Mapped[int] = mapped_column(Integer, primary_key=True)
|
||||
name: Mapped[str] = mapped_column(String(120), nullable=False)
|
||||
enabled: Mapped[bool] = mapped_column(Boolean, nullable=False, default=True)
|
||||
health_state: Mapped[HealthState] = mapped_column(Enum(HealthState), nullable=False, default=HealthState.HEALTHY)
|
||||
mywhoosh_state: Mapped[str] = mapped_column(String(32), nullable=False, default="unknown")
|
||||
garmin_state: Mapped[str] = mapped_column(String(32), nullable=False, default="unknown")
|
||||
action_reason: Mapped[str | None] = mapped_column(Text)
|
||||
mywhoosh_email_enc: Mapped[str] = mapped_column(Text, nullable=False)
|
||||
mywhoosh_password_enc: Mapped[str] = mapped_column(Text, nullable=False)
|
||||
garmin_email_enc: Mapped[str] = mapped_column(Text, nullable=False)
|
||||
garmin_password_enc: Mapped[str] = mapped_column(Text, nullable=False)
|
||||
created_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), default=utcnow)
|
||||
updated_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), default=utcnow, onupdate=utcnow)
|
||||
|
||||
|
||||
class Activity(Base):
|
||||
__tablename__ = "activities"
|
||||
__table_args__ = (UniqueConstraint("user_id", "mywhoosh_activity_id", name="uq_activity_user_mywhoosh"),)
|
||||
|
||||
id: Mapped[int] = mapped_column(Integer, primary_key=True)
|
||||
user_id: Mapped[int] = mapped_column(ForeignKey("sync_users.id", ondelete="CASCADE"), nullable=False, index=True)
|
||||
mywhoosh_activity_id: Mapped[str] = mapped_column(String(255), nullable=False)
|
||||
activity_name: Mapped[str] = mapped_column(String(255), nullable=False)
|
||||
activity_timestamp: Mapped[datetime | None] = mapped_column(DateTime(timezone=True))
|
||||
source_fit_path: Mapped[str | None] = mapped_column(Text)
|
||||
converted_fit_path: Mapped[str | None] = mapped_column(Text)
|
||||
status: Mapped[ActivityStatus] = mapped_column(Enum(ActivityStatus), nullable=False, default=ActivityStatus.DISCOVERED)
|
||||
last_completed_stage: Mapped[ActivityStatus] = mapped_column(Enum(ActivityStatus), nullable=False, default=ActivityStatus.DISCOVERED)
|
||||
retryable: Mapped[bool] = mapped_column(Boolean, nullable=False, default=True)
|
||||
garmin_activity_id: Mapped[str | None] = mapped_column(String(255))
|
||||
last_error: Mapped[str | None] = mapped_column(Text)
|
||||
created_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), default=utcnow)
|
||||
updated_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), default=utcnow, onupdate=utcnow)
|
||||
|
||||
|
||||
class SyncRun(Base):
|
||||
__tablename__ = "sync_runs"
|
||||
|
||||
id: Mapped[int] = mapped_column(Integer, primary_key=True)
|
||||
user_id: Mapped[int] = mapped_column(ForeignKey("sync_users.id", ondelete="CASCADE"), nullable=False, index=True)
|
||||
started_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), default=utcnow)
|
||||
finished_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True))
|
||||
status: Mapped[SyncRunStatus] = mapped_column(Enum(SyncRunStatus), nullable=False, default=SyncRunStatus.RUNNING)
|
||||
discovered_count: Mapped[int] = mapped_column(Integer, nullable=False, default=0)
|
||||
imported_count: Mapped[int] = mapped_column(Integer, nullable=False, default=0)
|
||||
skipped_count: Mapped[int] = mapped_column(Integer, nullable=False, default=0)
|
||||
failed_count: Mapped[int] = mapped_column(Integer, nullable=False, default=0)
|
||||
summary_error: Mapped[str | None] = mapped_column(Text)
|
||||
69
app/db/repositories.py
Normal file
69
app/db/repositories.py
Normal file
@@ -0,0 +1,69 @@
|
||||
from datetime import datetime
|
||||
|
||||
from sqlalchemy import select
|
||||
from sqlalchemy.exc import IntegrityError
|
||||
from sqlalchemy.orm import Session
|
||||
|
||||
from app.db.models import Activity, ActivityStatus, HealthState, SyncUser
|
||||
|
||||
|
||||
class UserRepository:
|
||||
def __init__(self, session: Session) -> None:
|
||||
self.session = session
|
||||
|
||||
def create(self, **values) -> SyncUser:
|
||||
user = SyncUser(**values)
|
||||
self.session.add(user)
|
||||
self.session.commit()
|
||||
return user
|
||||
|
||||
def get(self, user_id: int) -> SyncUser | None:
|
||||
return self.session.get(SyncUser, user_id)
|
||||
|
||||
def list_enabled(self) -> list[SyncUser]:
|
||||
return list(self.session.scalars(select(SyncUser).where(SyncUser.enabled.is_(True)).order_by(SyncUser.id)))
|
||||
|
||||
|
||||
class ActivityRepository:
|
||||
def __init__(self, session: Session) -> None:
|
||||
self.session = session
|
||||
|
||||
def get_or_create_discovered(
|
||||
self,
|
||||
*,
|
||||
user_id: int,
|
||||
mywhoosh_activity_id: str,
|
||||
activity_name: str,
|
||||
activity_timestamp: datetime | None,
|
||||
) -> tuple[Activity, bool]:
|
||||
existing = self.session.scalar(
|
||||
select(Activity).where(
|
||||
Activity.user_id == user_id,
|
||||
Activity.mywhoosh_activity_id == mywhoosh_activity_id,
|
||||
)
|
||||
)
|
||||
if existing is not None:
|
||||
return existing, False
|
||||
activity = Activity(
|
||||
user_id=user_id,
|
||||
mywhoosh_activity_id=mywhoosh_activity_id,
|
||||
activity_name=activity_name,
|
||||
activity_timestamp=activity_timestamp,
|
||||
status=ActivityStatus.DISCOVERED,
|
||||
last_completed_stage=ActivityStatus.DISCOVERED,
|
||||
)
|
||||
self.session.add(activity)
|
||||
try:
|
||||
self.session.commit()
|
||||
except IntegrityError:
|
||||
self.session.rollback()
|
||||
existing = self.session.scalar(
|
||||
select(Activity).where(
|
||||
Activity.user_id == user_id,
|
||||
Activity.mywhoosh_activity_id == mywhoosh_activity_id,
|
||||
)
|
||||
)
|
||||
if existing is None:
|
||||
raise
|
||||
return existing, False
|
||||
return activity, True
|
||||
18
app/db/session.py
Normal file
18
app/db/session.py
Normal file
@@ -0,0 +1,18 @@
|
||||
from sqlalchemy import create_engine
|
||||
from sqlalchemy.engine import Engine
|
||||
from sqlalchemy.orm import Session, sessionmaker
|
||||
|
||||
from app.db.models import Base
|
||||
|
||||
|
||||
def create_db_engine(database_url: str) -> Engine:
|
||||
connect_args = {"check_same_thread": False} if database_url.startswith("sqlite") else {}
|
||||
return create_engine(database_url, connect_args=connect_args, future=True)
|
||||
|
||||
|
||||
def create_session_factory(engine: Engine) -> sessionmaker[Session]:
|
||||
return sessionmaker(bind=engine, autoflush=False, expire_on_commit=False)
|
||||
|
||||
|
||||
def initialize_schema(engine: Engine) -> None:
|
||||
Base.metadata.create_all(engine)
|
||||
@@ -1,6 +1,7 @@
|
||||
from fastapi import FastAPI
|
||||
|
||||
from app.config import Settings, get_settings
|
||||
from app.db.session import create_db_engine, create_session_factory, initialize_schema
|
||||
|
||||
|
||||
def create_app(settings: Settings | None = None) -> FastAPI:
|
||||
@@ -12,6 +13,11 @@ def create_app(settings: Settings | None = None) -> FastAPI:
|
||||
app = FastAPI(title="MyWhoosh Garmin Sync")
|
||||
app.state.settings = resolved
|
||||
|
||||
engine = create_db_engine(resolved.database_url)
|
||||
initialize_schema(engine)
|
||||
app.state.db_engine = engine
|
||||
app.state.session_factory = create_session_factory(engine)
|
||||
|
||||
@app.get("/healthz")
|
||||
def healthz() -> dict[str, str]:
|
||||
return {"status": "ok"}
|
||||
|
||||
@@ -0,0 +1,30 @@
|
||||
import pytest
|
||||
from sqlalchemy import create_engine
|
||||
from sqlalchemy.orm import Session, sessionmaker
|
||||
from sqlalchemy.pool import StaticPool
|
||||
|
||||
from app.db.models import Base
|
||||
from app.db.repositories import ActivityRepository, UserRepository
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def db_session() -> Session:
|
||||
engine = create_engine(
|
||||
"sqlite://",
|
||||
connect_args={"check_same_thread": False},
|
||||
poolclass=StaticPool,
|
||||
)
|
||||
Base.metadata.create_all(engine)
|
||||
factory = sessionmaker(bind=engine, expire_on_commit=False)
|
||||
with factory() as session:
|
||||
yield session
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def user_repository(db_session: Session) -> UserRepository:
|
||||
return UserRepository(db_session)
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def activity_repository(db_session: Session) -> ActivityRepository:
|
||||
return ActivityRepository(db_session)
|
||||
|
||||
54
tests/db/test_repositories.py
Normal file
54
tests/db/test_repositories.py
Normal file
@@ -0,0 +1,54 @@
|
||||
from app.db.models import ActivityStatus, HealthState
|
||||
|
||||
|
||||
def test_create_two_independent_users(db_session, user_repository) -> None:
|
||||
first = user_repository.create(
|
||||
name="Max",
|
||||
enabled=True,
|
||||
health_state=HealthState.HEALTHY,
|
||||
mywhoosh_email_enc="mw-1",
|
||||
mywhoosh_password_enc="mw-pw-1",
|
||||
garmin_email_enc="g-1",
|
||||
garmin_password_enc="g-pw-1",
|
||||
)
|
||||
second = user_repository.create(
|
||||
name="Anna",
|
||||
enabled=True,
|
||||
health_state=HealthState.HEALTHY,
|
||||
mywhoosh_email_enc="mw-2",
|
||||
mywhoosh_password_enc="mw-pw-2",
|
||||
garmin_email_enc="g-2",
|
||||
garmin_password_enc="g-pw-2",
|
||||
)
|
||||
|
||||
assert first.id != second.id
|
||||
assert {u.name for u in user_repository.list_enabled()} == {"Max", "Anna"}
|
||||
|
||||
|
||||
def test_activity_external_id_is_unique_per_user(user_repository, activity_repository) -> None:
|
||||
user = user_repository.create(
|
||||
name="Max",
|
||||
enabled=True,
|
||||
health_state=HealthState.HEALTHY,
|
||||
mywhoosh_email_enc="a",
|
||||
mywhoosh_password_enc="b",
|
||||
garmin_email_enc="c",
|
||||
garmin_password_enc="d",
|
||||
)
|
||||
created, inserted = activity_repository.get_or_create_discovered(
|
||||
user_id=user.id,
|
||||
mywhoosh_activity_id="mw-123",
|
||||
activity_name="Morning Ride",
|
||||
activity_timestamp=None,
|
||||
)
|
||||
same, inserted_again = activity_repository.get_or_create_discovered(
|
||||
user_id=user.id,
|
||||
mywhoosh_activity_id="mw-123",
|
||||
activity_name="Morning Ride",
|
||||
activity_timestamp=None,
|
||||
)
|
||||
|
||||
assert inserted is True
|
||||
assert inserted_again is False
|
||||
assert created.id == same.id
|
||||
assert same.status == ActivityStatus.DISCOVERED
|
||||
Reference in New Issue
Block a user