Files
rf4-spotter/apps/api/app/routers/public_data.py
T
ik 55e5810e65
CI / backend-and-migrations (push) Waiting to run
CI / astro-build (push) Waiting to run
CI / dependency-audit (push) Waiting to run
CI / compose-e2e (push) Waiting to run
feat: search official records
2026-09-21 20:21:12 +07:00

178 lines
6.3 KiB
Python

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)
))