from datetime import datetime, timedelta, timezone from fastapi import APIRouter, HTTPException, Query from sqlalchemy import func, or_, select from sqlalchemy.orm import joinedload from ..config import settings from ..dependencies import Db from ..models import ( CatchReport, Bait, CommunityImportRun, DataSource, ExternalObservation, Fish, OfficialRecordImport, SourceType, Waterbody, ) from ..schemas import ( ImportRunPublicOut, OfficialRecordOut, PaginatedOfficialRecordOut, PublicObservationOut, SourceStatusOut, ) from ..time_utils import aware router = APIRouter() @router.get("/api/v1/community-observations", response_model=list[PublicObservationOut]) def community_observations( db: Db, limit: int = Query(12, ge=1, le=50), offset: int = Query(0, ge=0), waterbody: str | None = None, fish: str | None = None, ) -> list[PublicObservationOut]: query = select(ExternalObservation).join(ExternalObservation.source).options( joinedload(ExternalObservation.source) ).where( ExternalObservation.catch_report_id.is_(None), ExternalObservation.status.not_in(["rejected", "withdrawn"]), DataSource.enabled.is_(True), ) if waterbody: query = query.join(ExternalObservation.waterbody).where(Waterbody.slug == waterbody) if fish: query = query.join(ExternalObservation.fish).where(Fish.slug == fish) items = list(db.scalars( query.order_by(ExternalObservation.last_seen_at.desc(), ExternalObservation.id.desc()) .offset(offset).limit(limit) )) result: list[PublicObservationOut] = [] for item in items: missing = [] if item.x is None or item.y is None: missing.append("координаты") if item.weight_g is None: missing.append("вес") result.append(PublicObservationOut( id=item.id, source_system=item.source_system, source_name=item.source.name, source_url=item.source_url, fish_name=item.fish_name, waterbody_name=item.waterbody_name, x=item.x, y=item.y, weight_g=item.weight_g, last_seen_at=item.last_seen_at, missing_fields=missing, quality="incomplete" if missing else "unverified", )) return result @router.get("/api/v1/source-status", response_model=list[SourceStatusOut]) def source_status(db: Db) -> list[SourceStatusOut]: now = datetime.now(timezone.utc) result = [] for source in db.scalars(select(DataSource).order_by(DataSource.name)): runs = list(db.scalars( select(CommunityImportRun) .where(CommunityImportRun.source_system == source.key) .order_by(CommunityImportRun.started_at.desc()).limit(20) )) latest = runs[0] if runs else None success = next((run for run in runs if run.status == "success"), None) if not source.enabled: state = "disabled" elif latest is None: state = "waiting" elif latest.status == "failed": state = "source_changed" if "CommunityParseError" in (latest.error_summary or "") else "temporarily_limited" elif aware(latest.started_at) < now - timedelta(seconds=settings.community_import_interval_seconds * 2): state = "stale" else: state = "healthy" result.append(SourceStatusOut( source_system=source.key, name=source.name, status=state, last_started_at=latest.started_at if latest else None, last_success_at=success.started_at if success else None, observations=db.scalar( select(func.count()).select_from(ExternalObservation) .where(ExternalObservation.source_system == source.key) ) or 0, )) return result @router.get("/api/v1/records", response_model=PaginatedOfficialRecordOut) def records( db: Db, fish: str | None = None, waterbody: str | None = None, category: str | None = None, q: str | None = None, limit: int = Query(50, ge=1, le=100), offset: int = Query(0, ge=0), ) -> PaginatedOfficialRecordOut: if q is not None and len(q) > 100: raise HTTPException(status_code=422, detail="record search query is too long") query = select(CatchReport).options( joinedload(CatchReport.fish), joinedload(CatchReport.waterbody), joinedload(CatchReport.bait), ).where(CatchReport.source_type == SourceType.official_record) if fish: query = query.join(CatchReport.fish).where(Fish.slug == fish) if waterbody: query = query.join(CatchReport.waterbody).where(Waterbody.slug == waterbody) if category: query = query.where(CatchReport.raw_payload["category"].as_string() == category) if q and q.strip(): needle = f"%{q.strip()}%" query = query.outerjoin(Bait, CatchReport.bait_id == Bait.id).where(or_(CatchReport.player_name.ilike(needle), Bait.name.ilike(needle))) total = db.scalar( query.with_only_columns(func.count(CatchReport.id), maintain_column_froms=True).order_by(None) ) or 0 items = list(db.scalars( query.order_by(CatchReport.caught_at.desc(), CatchReport.weight_g.desc(), CatchReport.id.desc()) .offset(offset).limit(limit) )) return PaginatedOfficialRecordOut( items=[OfficialRecordOut( id=item.id, fish=item.fish.name_ru, weight_g=item.weight_g, waterbody=item.waterbody.name_ru, bait=item.bait.name if item.bait else None, player_name=item.player_name, record_date=item.caught_at, category=(item.raw_payload or {}).get("category"), region=(item.raw_payload or {}).get("region"), source_url=item.source_url, ) for item in items], total=total, limit=limit, offset=offset, ) @router.get("/api/v1/imports", response_model=list[ImportRunPublicOut]) def imports( db: Db, limit: int = Query(20, ge=1, le=100), offset: int = Query(0, ge=0), ) -> list[OfficialRecordImport]: return list(db.scalars( select(OfficialRecordImport) .order_by(OfficialRecordImport.started_at.desc(), OfficialRecordImport.id.desc()) .offset(offset).limit(limit) ))