from __future__ import annotations from datetime import datetime, timedelta, timezone import hashlib import hmac import json import logging import secrets import time as time_module from typing import Annotated, Literal from uuid import UUID import uuid import httpx from fastapi import Depends, FastAPI, File, Header, HTTPException, Query, Request, Response, UploadFile from fastapi.middleware.cors import CORSMiddleware from fastapi.responses import JSONResponse from sqlalchemy import case, func, or_, select from sqlalchemy.exc import IntegrityError from sqlalchemy.orm import Session, joinedload from .admin_security import verify_admin from .config import settings from .dependencies import Db from .community_review import ExternalReviewError, map_observation, publish_observation, reject_observation, suggest_aliases from .importer import ImportAlreadyRunning, ImportSourceError, import_records, normalize from .logging_config import configure_logging from .models import Bait, BaitKind, CatchReport, CommunityImportRun, DataSource, ExternalObservation, Fish, ModerationEvent, ModerationStatus, OfficialRecordImport, SourceType, Spot, SubmissionAttempt, Waterbody from .readiness import readiness_report from .routers.activity import router as activity_router from .routers.catalog import router as catalog_router from .routers.public_data import router as public_data_router from .routers.submissions import router as submissions_router from .time_utils import aware from .public_cache import public_cache from .schemas import ActivityOut, AdminCatchReportOut, AdminModerationHistoryOut, CatchReportAccepted, CatchReportCreate, CatchReportCreated, ExternalAliasSuggestionOut, ExternalObservationDecision, ExternalObservationMapping, ExternalObservationOut, ExternalObservationPublished, ImportRunOut, ModerationUpdate from .storage import ScreenshotError, client as storage_client, delete_screenshot, signed_screenshot_url, upload_screenshot from .submission_security import check_rate_limit from .submission_security import is_trusted_proxy as _is_trusted_proxy configure_logging(settings.log_level) logger = logging.getLogger("rf4.api") app = FastAPI(title="RF4 Spotter API", version="0.1.0") app.add_middleware( CORSMiddleware, allow_origins=settings.cors_origins, allow_methods=["GET", "POST", "PATCH", "DELETE"], allow_headers=["Authorization", "Content-Type"], ) @app.middleware("http") async def structured_request_log(request: Request, call_next): request_id = uuid.uuid4().hex started = time_module.perf_counter() status_code = 500 try: response = await call_next(request) status_code = response.status_code response.headers["X-Request-ID"] = request_id response.headers["X-Content-Type-Options"] = "nosniff" response.headers["Referrer-Policy"] = "strict-origin-when-cross-origin" response.headers["Permissions-Policy"] = "camera=(), microphone=(), geolocation=()" response.headers["X-Frame-Options"] = "DENY" response.headers["Cross-Origin-Opener-Policy"] = "same-origin" if request.url.path.startswith("/api/v1/admin/") or request.url.path == "/api/v1/catch-reports": response.headers["Cache-Control"] = "no-store" if settings.deployment_environment == "production": response.headers["Strict-Transport-Security"] = "max-age=31536000; includeSubDomains" return response except Exception as exc: logger.error("request failed", extra={"request_id": request_id, "error_type": type(exc).__name__}) raise finally: logger.log( logging.DEBUG if request.url.path in {"/health", "/ready"} else logging.INFO, "request completed", extra={ "request_id": request_id, "method": request.method, "path": request.url.path, "status_code": status_code, "duration_ms": round((time_module.perf_counter() - started) * 1000, 2), }, ) @app.get("/health") def health() -> dict[str, str]: return {"status": "ok"} @app.get("/ready") def ready(db: Db) -> JSONResponse: is_ready, components = readiness_report( db, storage_client(), import_required=settings.official_import_required, import_interval_seconds=settings.import_interval_seconds, community_import_interval_seconds=settings.community_import_interval_seconds, ) return JSONResponse( status_code=200 if is_ready else 503, content={"status": "ready" if is_ready else "not_ready", "version": settings.app_version, "revision": settings.app_revision, "components": components}, ) app.include_router(catalog_router) app.include_router(activity_router) app.include_router(public_data_router) def _admin(request: Request, db: Db, authorization: Annotated[str | None, Header()] = None) -> str: return verify_admin(request, db, authorization, settings) @app.get("/api/v1/admin/diagnostics") def admin_diagnostics(db: Db, _: Annotated[str, Depends(_admin)]) -> JSONResponse: report_counts = {status.value: count for status, count in db.execute( select(CatchReport.moderation_status, func.count()).group_by(CatchReport.moderation_status) )} observation_counts = {status: count for status, count in db.execute( select(ExternalObservation.status, func.count()).group_by(ExternalObservation.status) )} payload = { "generated_at": datetime.now(timezone.utc).isoformat(), "build": {"version": settings.app_version, "revision": settings.app_revision, "environment": settings.deployment_environment}, "counts": { "catch_reports": report_counts, "external_observations": observation_counts, "data_sources": db.scalar(select(func.count()).select_from(DataSource)) or 0, "enabled_data_sources": db.scalar(select(func.count()).select_from(DataSource).where(DataSource.enabled.is_(True))) or 0, "official_import_runs": db.scalar(select(func.count()).select_from(OfficialRecordImport)) or 0, "community_import_runs": db.scalar(select(func.count()).select_from(CommunityImportRun)) or 0, }, } return JSONResponse(payload, headers={"Content-Disposition": "attachment; filename=rf4spotter-diagnostics.json"}) @app.get("/api/v1/admin/moderation-history", response_model=list[AdminModerationHistoryOut]) def admin_moderation_history( db: Db, _: Annotated[str, Depends(_admin)], limit: int = Query(50, ge=1, le=200), offset: int = Query(0, ge=0), ) -> list[AdminModerationHistoryOut]: report_events = list(db.scalars( select(ModerationEvent).order_by(ModerationEvent.created_at.desc()).limit(limit + offset) )) external_events = list(db.scalars( select(ExternalObservation).where(ExternalObservation.reviewed_at.is_not(None)) .order_by(ExternalObservation.reviewed_at.desc()).limit(limit + offset) )) history = [AdminModerationHistoryOut( entity_type="catch_report", entity_id=event.catch_report_id, decided_at=event.created_at, action=event.new_status.value, moderator=event.moderator, reason=event.reason, ) for event in report_events] history.extend(AdminModerationHistoryOut( entity_type="external_observation", entity_id=observation.id, decided_at=observation.reviewed_at, action=observation.status, moderator=None, reason=observation.review_note, ) for observation in external_events if observation.reviewed_at is not None) history.sort(key=lambda event: aware(event.decided_at), reverse=True) return history[offset:offset + limit] @app.get("/api/v1/admin/moderation-history-export") def admin_moderation_history_export( db: Db, _: Annotated[str, Depends(_admin)], limit: int = Query(1000, ge=1, le=5000), ) -> JSONResponse: """Return an anonymized, analysis-safe decision export.""" events = admin_moderation_history(db, _, limit=limit, offset=0) payload = { "generated_at": datetime.now(timezone.utc).isoformat(), "count": len(events), "events": [{ "entity_type": event.entity_type, "decided_at": event.decided_at.isoformat(), "action": event.action, "requires_confirmation": event.requires_confirmation, } for event in events], } return JSONResponse(payload, headers={ "Content-Disposition": "attachment; filename=rf4spotter-moderation-history.json", }) @app.get("/api/v1/admin/imports", response_model=list[ImportRunOut]) def admin_imports( db: Db, _: Annotated[str, Depends(_admin)], limit: int = Query(20, ge=1, le=100), offset: int = Query(0, ge=0), ) -> list[OfficialRecordImport]: query = select(OfficialRecordImport).order_by( OfficialRecordImport.started_at.desc(), OfficialRecordImport.id.desc() ).offset(offset).limit(limit) return list(db.scalars(query)) @app.post("/api/v1/admin/imports/official-records", response_model=ImportRunOut, status_code=201) def admin_start_official_import(db: Db, _: Annotated[str, Depends(_admin)]) -> OfficialRecordImport: try: return import_records( db, url=settings.official_records_url, region=settings.official_records_region, category=settings.official_records_category, ) except ImportAlreadyRunning as exc: raise HTTPException(status_code=409, detail=str(exc)) from exc except (ImportSourceError, httpx.HTTPError) as exc: raise HTTPException(status_code=502, detail=f"official records import failed: {exc}") from exc def _external_out(item: ExternalObservation) -> ExternalObservationOut: allowed_payload = { key: value for key, value in (item.payload or {}).items() if key in { "bait", "fishing_method", "rig_type", "retrieve_method", "retrieve_speed", "player_name", "published_at", "region", "category", } and (value is None or isinstance(value, (str, int, float, bool))) } missing_fields = [] if item.x is None or item.y is None: missing_fields.append("coordinates") if item.weight_g is None: missing_fields.append("weight_g") return ExternalObservationOut( id=item.id, source_system=item.source_system, source_external_id=item.source_external_id, source_url=item.source_url, fish_name=item.fish_name, fish_external_id=item.fish_external_id, waterbody_name=item.waterbody_name, waterbody_external_id=item.waterbody_external_id, x=item.x, y=item.y, weight_g=item.weight_g, published_at=item.published_at, first_seen_at=item.first_seen_at, last_seen_at=item.last_seen_at, reviewed_at=item.reviewed_at, status=item.status, fish_slug=item.fish.slug if item.fish else None, waterbody_slug=item.waterbody.slug if item.waterbody else None, catch_report_id=item.catch_report_id, review_note=item.review_note, missing_fields=missing_fields, source_payload=allowed_payload, ) @app.get("/api/v1/admin/external-observations", response_model=list[ExternalObservationOut]) def admin_external_observations( db: Db, _: Annotated[str, Depends(_admin)], status: Literal["staged", "mapped", "ready", "published", "rejected", "review"] | None = None, source_system: str | None = None, completeness: Literal["all", "complete", "incomplete"] = "all", order: Literal["newest", "oldest", "risk"] = "newest", q: str | None = Query(None, max_length=100), limit: int = Query(50, ge=1, le=200), offset: int = Query(0, ge=0), ) -> list[ExternalObservationOut]: query = select(ExternalObservation).options( joinedload(ExternalObservation.fish), joinedload(ExternalObservation.waterbody), ) if status == "review": query = query.where(ExternalObservation.status.in_(["staged", "mapped", "ready"])) elif status: query = query.where(ExternalObservation.status == status) if source_system: query = query.where(ExternalObservation.source_system == source_system) if completeness == "complete": query = query.where( ExternalObservation.x.is_not(None), ExternalObservation.y.is_not(None), ExternalObservation.weight_g.is_not(None), ) elif completeness == "incomplete": query = query.where(or_( ExternalObservation.x.is_(None), ExternalObservation.y.is_(None), ExternalObservation.weight_g.is_(None), )) if q and q.strip(): term = q.strip() query = query.where(or_( ExternalObservation.fish_name.icontains(term, autoescape=True), ExternalObservation.waterbody_name.icontains(term, autoescape=True), )) if order == "risk": incomplete = case( (or_(ExternalObservation.x.is_(None), ExternalObservation.y.is_(None), ExternalObservation.weight_g.is_(None)), 0), else_=1, ) workflow = case( (ExternalObservation.status == "staged", 0), (ExternalObservation.status == "mapped", 1), else_=2, ) ordering = (incomplete, workflow, ExternalObservation.last_seen_at.asc(), ExternalObservation.id.desc()) else: direction = ExternalObservation.last_seen_at.asc() if order == "oldest" else ExternalObservation.last_seen_at.desc() ordering = (direction, ExternalObservation.id.desc()) items = db.scalars(query.order_by(*ordering).offset(offset).limit(limit)) return [_external_out(item) for item in items] @app.get("/api/v1/admin/external-observations/{observation_id}/alias-suggestions", response_model=ExternalAliasSuggestionOut) def admin_external_alias_suggestions( observation_id: UUID, db: Db, _: Annotated[str, Depends(_admin)], ) -> ExternalAliasSuggestionOut: observation = db.get(ExternalObservation, observation_id) if observation is None: raise HTTPException(status_code=404, detail="external observation not found") fish, waterbody = suggest_aliases(db, observation) return ExternalAliasSuggestionOut( fish_slug=fish.slug if fish else None, waterbody_slug=waterbody.slug if waterbody else None, ) @app.patch("/api/v1/admin/external-observations/{observation_id}/mapping", response_model=ExternalObservationOut) def admin_map_external_observation( observation_id: UUID, payload: ExternalObservationMapping, db: Db, _: Annotated[str, Depends(_admin)], ) -> ExternalObservationOut: observation = db.get(ExternalObservation, observation_id) fish = db.scalar(select(Fish).where(Fish.slug == payload.fish_slug)) waterbody = db.scalar(select(Waterbody).where(Waterbody.slug == payload.waterbody_slug)) if observation is None: raise HTTPException(status_code=404, detail="external observation not found") if fish is None or waterbody is None: raise HTTPException(status_code=422, detail="unknown fish or waterbody") try: return _external_out(map_observation(db, observation, fish, waterbody, note=payload.note)) except ExternalReviewError as exc: raise HTTPException(status_code=409, detail=str(exc)) from exc @app.post("/api/v1/admin/external-observations/{observation_id}/publish", response_model=ExternalObservationPublished) def admin_publish_external_observation( observation_id: UUID, db: Db, _: Annotated[str, Depends(_admin)], ) -> ExternalObservationPublished: observation = db.get(ExternalObservation, observation_id) if observation is None: raise HTTPException(status_code=404, detail="external observation not found") try: report = publish_observation(db, observation) except ExternalReviewError as exc: raise HTTPException(status_code=409, detail=str(exc)) from exc public_cache.invalidate() return ExternalObservationPublished(observation_id=observation.id, catch_report_id=report.id, status=observation.status) @app.patch("/api/v1/admin/external-observations/{observation_id}/reject", response_model=ExternalObservationOut) def admin_reject_external_observation( observation_id: UUID, payload: ExternalObservationDecision, db: Db, _: Annotated[str, Depends(_admin)], ) -> ExternalObservationOut: observation = db.get(ExternalObservation, observation_id) if observation is None: raise HTTPException(status_code=404, detail="external observation not found") try: return _external_out(reject_observation(db, observation, reason=payload.reason)) except ExternalReviewError as exc: raise HTTPException(status_code=409, detail=str(exc)) from exc @submissions_router.post("/api/v1/catch-reports", response_model=CatchReportAccepted, status_code=201) def create_catch_report( payload: CatchReportCreate, request: Request, db: Db, idempotency_key: Annotated[str | None, Header()] = None, ) -> CatchReportAccepted: if payload.website: raise HTTPException(status_code=400, detail="invalid submission") payload_hash = hashlib.sha256(json.dumps(payload.model_dump(mode="json"), sort_keys=True, separators=(",", ":")).encode()).hexdigest() # A05: Server-side idempotency — check BEFORE rate limit to avoid polluting table if idempotency_key: key_hash = hmac.new(settings.rate_limit_secret.encode(), idempotency_key.encode(), hashlib.sha256).hexdigest() cutoff = datetime.now(timezone.utc) - timedelta(minutes=5) # Force refresh from database to see committed data from previous requests db.expire_all() existing = db.scalar( select(SubmissionAttempt).where( SubmissionAttempt.idempotency_key == key_hash, SubmissionAttempt.created_at >= cutoff, ) ) if existing is not None: # Return 200 with idempotent flag — client can retry safely logger.info("idempotent hit", extra={"idempotency_key": idempotency_key[:8]}) report = existing.catch_report if existing.payload_hash and not hmac.compare_digest(existing.payload_hash, payload_hash): raise HTTPException(status_code=409, detail="Idempotency-Key was already used with different payload") if report is None: raise HTTPException(status_code=409, detail="idempotency record is incomplete; retry with a new key") # Re-derive the one-time upload token from the idempotency key; # only its hash is persisted, so the secret is never stored. replay_token = hmac.new(settings.rate_limit_secret.encode(), (key_hash + ":upload").encode(), hashlib.sha256).hexdigest() if not hmac.compare_digest(hashlib.sha256(replay_token.encode()).hexdigest(), report.screenshot_upload_token_hash or ""): raise HTTPException(status_code=409, detail="idempotency record token mismatch; retry with a new key") return JSONResponse(status_code=200, content={"id": str(report.id), "moderation_status": report.moderation_status.value, "screenshot_upload_token": replay_token, "idempotent": True}) logger.info("idempotency check miss", extra={"idempotency_key": idempotency_key[:8]}) _check_rate_limit(request, db) fish = db.scalar(select(Fish).where(Fish.slug == payload.fish_slug)) waterbody = db.scalar(select(Waterbody).where(Waterbody.slug == payload.waterbody_slug)) if fish is None or waterbody is None: raise HTTPException(status_code=422, detail="unknown fish or waterbody") spot = db.scalar(select(Spot).where(Spot.waterbody_id == waterbody.id, Spot.x == payload.x, Spot.y == payload.y)) if spot is None: spot = Spot(waterbody=waterbody, x=payload.x, y=payload.y) db.add(spot) bait = None if payload.bait_name and payload.bait_name.strip(): key = normalize(payload.bait_name) bait = db.scalar(select(Bait).where(Bait.normalized_name == key)) if bait is None: bait = Bait(name=payload.bait_name.strip(), normalized_name=key, kind=BaitKind.unknown) db.add(bait) upload_token = (hmac.new(settings.rate_limit_secret.encode(), (key_hash + ":upload").encode(), hashlib.sha256).hexdigest() if idempotency_key else secrets.token_urlsafe(32)) report = CatchReport(fish=fish, spot=spot, waterbody=waterbody, bait=bait, weight_g=payload.weight_g, fishing_method=payload.fishing_method, rig_type=payload.rig_type, retrieve_method=payload.retrieve_method, retrieve_speed=payload.retrieve_speed, caught_at=payload.caught_at, reported_at=datetime.now(timezone.utc), player_name=payload.player_name, source_type=SourceType.user, source_url=payload.source_url, source_confidence=60, moderation_status=ModerationStatus.pending, raw_payload={"comment": payload.comment} if payload.comment else None, screenshot_upload_token_hash=hashlib.sha256(upload_token.encode()).hexdigest()) db.add(report) # Store idempotency key if provided if idempotency_key: key_hash = hmac.new(settings.rate_limit_secret.encode(), idempotency_key.encode(), hashlib.sha256).hexdigest() # One transaction: a unique-key race must roll back the report too. db.add(SubmissionAttempt(client_hash="", idempotency_key=key_hash, catch_report=report, payload_hash=payload_hash, created_at=datetime.now(timezone.utc))) try: db.commit() except IntegrityError: # Another request won the same idempotency key race. db.rollback() winner = db.scalar(select(SubmissionAttempt).where(SubmissionAttempt.idempotency_key == key_hash)) if winner and winner.catch_report: replay_token = hmac.new(settings.rate_limit_secret.encode(), (key_hash + ":upload").encode(), hashlib.sha256).hexdigest() return JSONResponse(status_code=200, content={"id": str(winner.catch_report.id), "moderation_status": winner.catch_report.moderation_status.value, "screenshot_upload_token": replay_token, "idempotent": True}) raise logger.info("idempotency key stored", extra={"idempotency_key": idempotency_key[:8]}) else: db.commit() return CatchReportAccepted(id=report.id, moderation_status=report.moderation_status.value, screenshot_upload_token=upload_token, idempotent=False) @submissions_router.post("/api/v1/catch-reports/{report_id}/screenshot", status_code=204, response_class=Response) def add_screenshot( report_id: UUID, db: Db, screenshot: UploadFile = File(), upload_token: Annotated[str | None, Header(alias="X-Upload-Token")] = None, ) -> Response: report = db.get(CatchReport, report_id) if report is None or report.source_type != SourceType.user or report.moderation_status != ModerationStatus.pending: raise HTTPException(status_code=404, detail="pending catch report not found") supplied_hash = hashlib.sha256((upload_token or "").encode()).hexdigest() if not report.screenshot_upload_token_hash or not hmac.compare_digest(report.screenshot_upload_token_hash, supplied_hash): raise HTTPException(status_code=401, detail="invalid screenshot upload token") if report.screenshot_key: raise HTTPException(status_code=409, detail="screenshot already uploaded") raw = screenshot.file.read(settings.screenshot_max_bytes + 1) try: report.screenshot_key = upload_screenshot(raw, filename=screenshot.filename, content_type=screenshot.content_type) except ScreenshotError as exc: raise HTTPException(status_code=422, detail=str(exc)) from exc report.screenshot_upload_token_hash = None db.commit() return Response(status_code=204) app.include_router(submissions_router) @app.get("/api/v1/admin/catch-reports", response_model=list[AdminCatchReportOut]) def admin_reports(db: Db, _: Annotated[str, Depends(_admin)], status: ModerationStatus = ModerationStatus.pending, limit: int = Query(50, ge=1, le=100), offset: int = Query(0, ge=0)) -> list[AdminCatchReportOut]: reports = list(db.scalars(select(CatchReport).options(joinedload(CatchReport.fish), joinedload(CatchReport.waterbody), joinedload(CatchReport.spot), joinedload(CatchReport.bait)).where(CatchReport.source_type == SourceType.user, CatchReport.moderation_status == status, CatchReport.deleted_at.is_(None)).order_by(CatchReport.reported_at, CatchReport.id).offset(offset).limit(limit))) return [AdminCatchReportOut(id=r.id, fish=r.fish.name_ru, waterbody=r.waterbody.name_ru, coordinates=f"{r.spot.x}:{r.spot.y}" if r.spot else "—", weight_g=r.weight_g, bait=r.bait.name if r.bait else None, player_name=r.player_name, reported_at=r.reported_at, moderation_status=r.moderation_status.value, comment=(r.raw_payload or {}).get("comment"), screenshot_url=signed_screenshot_url(r.screenshot_key) if r.screenshot_key else None) for r in reports] @app.patch("/api/v1/admin/catch-reports/{report_id}", response_model=CatchReportCreated) def moderate_report(report_id: UUID, payload: ModerationUpdate, db: Db, moderator: Annotated[str, Depends(_admin)]) -> CatchReportCreated: report = db.get(CatchReport, report_id) if report is None or report.source_type != SourceType.user or report.deleted_at is not None: raise HTTPException(status_code=404, detail="catch report not found") previous = report.moderation_status report.moderation_status = ModerationStatus(payload.status) db.add(ModerationEvent(catch_report=report, created_at=datetime.now(timezone.utc), previous_status=previous, new_status=report.moderation_status, moderator=moderator, reason=payload.reason)) db.commit() public_cache.invalidate() return CatchReportCreated(id=report.id, moderation_status=report.moderation_status.value) @app.delete("/api/v1/admin/catch-reports/{report_id}", status_code=204, response_class=Response) def delete_report(report_id: UUID, db: Db, moderator: Annotated[str, Depends(_admin)]) -> Response: report = db.get(CatchReport, report_id) if report is None or report.source_type != SourceType.user or report.deleted_at is not None: raise HTTPException(status_code=404, detail="catch report not found") previous = report.moderation_status if report.screenshot_key: try: delete_screenshot(report.screenshot_key) except Exception as exc: raise HTTPException(status_code=502, detail="screenshot deletion failed") from exc report.moderation_status = ModerationStatus.rejected report.deleted_at = datetime.now(timezone.utc) report.player_name = None report.source_url = None report.screenshot_key = None report.raw_payload = None db.add(ModerationEvent(catch_report=report, created_at=report.deleted_at, previous_status=previous, new_status=ModerationStatus.rejected, moderator=moderator, reason="user report deleted and anonymized")) db.commit() public_cache.invalidate() return Response(status_code=204) def _check_rate_limit(request: Request, db: Session) -> None: check_rate_limit(request, db, settings)