feat: add external observation review queue

This commit is contained in:
ik
2026-09-03 18:47:31 +07:00
parent 53d700126a
commit c4f4deb24b
10 changed files with 429 additions and 13 deletions
@@ -0,0 +1,45 @@
"""Add canonical aliases and external-observation review state."""
from alembic import op
import sqlalchemy as sa
revision = "0009"
down_revision = "0008"
branch_labels = None
depends_on = None
def upgrade() -> None:
op.add_column("external_observation", sa.Column("fish_id", sa.Uuid(), sa.ForeignKey("fish.id")))
op.add_column("external_observation", sa.Column("waterbody_id", sa.Uuid(), sa.ForeignKey("waterbody.id")))
op.add_column("external_observation", sa.Column("catch_report_id", sa.Uuid(), sa.ForeignKey("catch_report.id")))
op.add_column("external_observation", sa.Column("review_note", sa.Text()))
op.add_column("external_observation", sa.Column("reviewed_at", sa.DateTime(timezone=True)))
op.create_unique_constraint("uq_external_observation_catch_report_id", "external_observation", ["catch_report_id"])
op.create_table(
"external_entity_alias",
sa.Column("id", sa.Uuid(), primary_key=True),
sa.Column("source_system", sa.String(50), sa.ForeignKey("data_source.key"), nullable=False),
sa.Column("entity_type", sa.String(20), nullable=False),
sa.Column("external_id", sa.String(200), nullable=False),
sa.Column("external_name", sa.String(200), nullable=False),
sa.Column("fish_id", sa.Uuid(), sa.ForeignKey("fish.id")),
sa.Column("waterbody_id", sa.Uuid(), sa.ForeignKey("waterbody.id")),
sa.Column("updated_at", sa.DateTime(timezone=True), nullable=False),
sa.CheckConstraint(
"(entity_type = 'fish' AND fish_id IS NOT NULL AND waterbody_id IS NULL) OR "
"(entity_type = 'waterbody' AND waterbody_id IS NOT NULL AND fish_id IS NULL)",
name="ck_external_entity_alias_target",
),
sa.UniqueConstraint("source_system", "entity_type", "external_id"),
)
op.create_index("ix_external_entity_alias_source_system", "external_entity_alias", ["source_system"])
def downgrade() -> None:
op.drop_table("external_entity_alias")
op.drop_constraint("uq_external_observation_catch_report_id", "external_observation", type_="unique")
op.drop_column("external_observation", "reviewed_at")
op.drop_column("external_observation", "review_note")
op.drop_column("external_observation", "catch_report_id")
op.drop_column("external_observation", "waterbody_id")
op.drop_column("external_observation", "fish_id")
+131
View File
@@ -0,0 +1,131 @@
from __future__ import annotations
import hashlib
from datetime import datetime, timezone
from sqlalchemy import select
from sqlalchemy.orm import Session
from .importer import normalize
from .models import (
Bait, BaitKind, CatchReport, ExternalEntityAlias, ExternalObservation,
Fish, ModerationStatus, SourceType, Spot, Waterbody,
)
class ExternalReviewError(ValueError):
pass
def map_observation(
session: Session, observation: ExternalObservation, fish: Fish, waterbody: Waterbody,
*, note: str | None = None,
) -> ExternalObservation:
if observation.status == "published":
raise ExternalReviewError("published observation cannot be remapped")
observation.fish = fish
observation.waterbody = waterbody
observation.review_note = note
observation.reviewed_at = datetime.now(timezone.utc)
observation.status = "ready" if _complete(observation) else "mapped"
_save_alias(session, observation, "fish", observation.fish_external_id or observation.fish_name, fish=fish)
_save_alias(session, observation, "waterbody", observation.waterbody_external_id or observation.waterbody_name, waterbody=waterbody)
session.commit()
return observation
def reject_observation(session: Session, observation: ExternalObservation, *, reason: str) -> ExternalObservation:
if observation.status == "published":
raise ExternalReviewError("published observation cannot be rejected")
observation.status = "rejected"
observation.review_note = reason
observation.reviewed_at = datetime.now(timezone.utc)
session.commit()
return observation
def publish_observation(session: Session, observation: ExternalObservation) -> CatchReport:
if observation.catch_report is not None:
return observation.catch_report
if observation.fish is None or observation.waterbody is None or not _complete(observation):
raise ExternalReviewError("fish, waterbody, coordinates and weight are required for publication")
spot = session.scalar(select(Spot).where(
Spot.waterbody_id == observation.waterbody.id, Spot.x == observation.x, Spot.y == observation.y,
))
if spot is None:
spot = Spot(waterbody=observation.waterbody, x=observation.x, y=observation.y)
session.add(spot)
bait = _bait(session, observation.payload.get("bait"))
now = datetime.now(timezone.utc)
report = CatchReport(
fish=observation.fish, waterbody=observation.waterbody, spot=spot, bait=bait,
weight_g=observation.weight_g,
fishing_method=observation.payload.get("fishing_method"),
rig_type=observation.payload.get("rig_type"),
retrieve_method=observation.payload.get("retrieve_method"),
retrieve_speed=observation.payload.get("retrieve_speed"),
caught_at=observation.published_at,
reported_at=observation.published_at or observation.first_seen_at,
player_name=observation.payload.get("player_name"),
source_type=SourceType.manual_import,
source_url=observation.source_url,
source_external_id=_report_external_id(observation),
source_confidence=observation.source.default_confidence,
moderation_status=ModerationStatus.approved,
raw_payload={
"provenance": {
"external_observation_id": str(observation.id),
"source_system": observation.source_system,
"source_external_id": observation.source_external_id,
},
"original": observation.payload,
},
)
session.add(report)
session.flush()
observation.catch_report = report
observation.status = "published"
observation.reviewed_at = now
session.commit()
return report
def _complete(observation: ExternalObservation) -> bool:
return observation.x is not None and observation.y is not None and observation.weight_g is not None
def _save_alias(
session: Session, observation: ExternalObservation, entity_type: str, external_id: str,
*, fish: Fish | None = None, waterbody: Waterbody | None = None,
) -> None:
alias = session.scalar(select(ExternalEntityAlias).where(
ExternalEntityAlias.source_system == observation.source_system,
ExternalEntityAlias.entity_type == entity_type,
ExternalEntityAlias.external_id == external_id,
))
if alias is None:
alias = ExternalEntityAlias(
source_system=observation.source_system, entity_type=entity_type,
external_id=external_id, external_name=observation.fish_name if fish else observation.waterbody_name,
)
session.add(alias)
alias.fish = fish
alias.waterbody = waterbody
alias.updated_at = datetime.now(timezone.utc)
def _bait(session: Session, value: object) -> Bait | None:
name = str(value or "").strip()
if not name:
return None
key = normalize(name)
bait = session.scalar(select(Bait).where(Bait.normalized_name == key))
if bait is None:
bait = Bait(name=name[:200], normalized_name=key[:200], kind=BaitKind.unknown)
session.add(bait)
return bait
def _report_external_id(observation: ExternalObservation) -> str:
raw = f"{observation.source_system}:{observation.source_external_id}".encode()
return "ext:" + hashlib.sha256(raw).hexdigest()[:60]
+79 -3
View File
@@ -16,9 +16,10 @@ from sqlalchemy.orm import Session, joinedload
from .activity import activity_rows
from .database import get_session
from .config import settings
from .community_review import ExternalReviewError, map_observation, publish_observation, reject_observation
from .importer import ImportSourceError, import_records, normalize
from .models import Bait, BaitKind, CatchReport, Fish, ModerationEvent, ModerationStatus, OfficialRecordImport, SourceType, Spot, SubmissionAttempt, Waterbody
from .schemas import ActivityOut, AdminCatchReportOut, BaitOut, CatchOut, CatchReportCreate, CatchReportCreated, FishOut, ImportRunOut, ModerationUpdate, OfficialRecordOut, SpotOut, WaterbodyOut
from .models import Bait, BaitKind, CatchReport, ExternalObservation, Fish, ModerationEvent, ModerationStatus, OfficialRecordImport, SourceType, Spot, SubmissionAttempt, Waterbody
from .schemas import ActivityOut, AdminCatchReportOut, BaitOut, CatchOut, CatchReportCreate, CatchReportCreated, ExternalObservationDecision, ExternalObservationMapping, ExternalObservationOut, ExternalObservationPublished, FishOut, ImportRunOut, ModerationUpdate, OfficialRecordOut, SpotOut, WaterbodyOut
from .storage import ScreenshotError, delete_screenshot, signed_screenshot_url, upload_screenshot
@@ -26,7 +27,7 @@ app = FastAPI(title="RF4 Spotter API", version="0.1.0")
app.add_middleware(
CORSMiddleware,
allow_origins=["http://localhost:4321", "http://127.0.0.1:4321"],
allow_methods=["GET", "PATCH", "DELETE"],
allow_methods=["GET", "POST", "PATCH", "DELETE"],
allow_headers=["Authorization", "Content-Type"],
)
Db = Annotated[Session, Depends(get_session)]
@@ -147,6 +148,81 @@ def admin_start_official_import(db: Db, _: Annotated[str, Depends(_admin)]) -> O
raise HTTPException(status_code=502, detail=f"official records import failed: {exc}") from exc
def _external_out(item: ExternalObservation) -> ExternalObservationOut:
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,
last_seen_at=item.last_seen_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,
)
@app.get("/api/v1/admin/external-observations", response_model=list[ExternalObservationOut])
def admin_external_observations(
db: Db, _: Annotated[str, Depends(_admin)], status: str | None = None,
source_system: str | None = None, 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:
query = query.where(ExternalObservation.status == status)
if source_system:
query = query.where(ExternalObservation.source_system == source_system)
items = db.scalars(query.order_by(ExternalObservation.last_seen_at.desc()).offset(offset).limit(limit))
return [_external_out(item) for item in items]
@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
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
@app.post("/api/v1/catch-reports", response_model=CatchReportCreated, status_code=201)
def create_catch_report(payload: CatchReportCreate, request: Request, db: Db) -> CatchReportCreated:
if payload.website:
+31 -1
View File
@@ -4,7 +4,7 @@ import enum
import uuid
from datetime import datetime, time
from sqlalchemy import JSON, DateTime, Enum, ForeignKey, Integer, String, Text, Time, UniqueConstraint
from sqlalchemy import JSON, CheckConstraint, DateTime, Enum, ForeignKey, Integer, String, Text, Time, UniqueConstraint
from sqlalchemy.orm import Mapped, mapped_column, relationship
from .database import Base
@@ -166,4 +166,34 @@ class ExternalObservation(Base):
last_seen_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), index=True)
status: Mapped[str] = mapped_column(String(30), default="staged", index=True)
payload: Mapped[dict] = mapped_column(JSON)
fish_id: Mapped[uuid.UUID | None] = mapped_column(ForeignKey("fish.id"))
waterbody_id: Mapped[uuid.UUID | None] = mapped_column(ForeignKey("waterbody.id"))
catch_report_id: Mapped[uuid.UUID | None] = mapped_column(ForeignKey("catch_report.id"), unique=True)
review_note: Mapped[str | None] = mapped_column(Text)
reviewed_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True))
source: Mapped[DataSource] = relationship()
fish: Mapped[Fish | None] = relationship()
waterbody: Mapped[Waterbody | None] = relationship()
catch_report: Mapped[CatchReport | None] = relationship()
class ExternalEntityAlias(Base):
__tablename__ = "external_entity_alias"
__table_args__ = (
UniqueConstraint("source_system", "entity_type", "external_id"),
CheckConstraint(
"(entity_type = 'fish' AND fish_id IS NOT NULL AND waterbody_id IS NULL) OR "
"(entity_type = 'waterbody' AND waterbody_id IS NOT NULL AND fish_id IS NULL)",
name="ck_external_entity_alias_target",
),
)
id: Mapped[uuid.UUID] = mapped_column(primary_key=True, default=uuid.uuid4)
source_system: Mapped[str] = mapped_column(ForeignKey("data_source.key"), index=True)
entity_type: Mapped[str] = mapped_column(String(20))
external_id: Mapped[str] = mapped_column(String(200))
external_name: Mapped[str] = mapped_column(String(200))
fish_id: Mapped[uuid.UUID | None] = mapped_column(ForeignKey("fish.id"))
waterbody_id: Mapped[uuid.UUID | None] = mapped_column(ForeignKey("waterbody.id"))
updated_at: Mapped[datetime] = mapped_column(DateTime(timezone=True))
fish: Mapped[Fish | None] = relationship()
waterbody: Mapped[Waterbody | None] = relationship()
+37
View File
@@ -161,3 +161,40 @@ class ModerationUpdate(BaseModel):
if value not in {"approved", "rejected", "pending"}:
raise ValueError("unknown moderation status")
return value
class ExternalObservationOut(BaseModel):
id: UUID
source_system: str
source_external_id: str
source_url: str
fish_name: str
fish_external_id: str | None
waterbody_name: str
waterbody_external_id: str | None
x: int | None
y: int | None
weight_g: int | None
published_at: datetime | None
last_seen_at: datetime
status: str
fish_slug: str | None
waterbody_slug: str | None
catch_report_id: UUID | None
review_note: str | None
class ExternalObservationMapping(BaseModel):
fish_slug: str
waterbody_slug: str
note: str | None = Field(default=None, max_length=1000)
class ExternalObservationDecision(BaseModel):
reason: str = Field(min_length=1, max_length=1000)
class ExternalObservationPublished(BaseModel):
observation_id: UUID
catch_report_id: UUID
status: str
+63 -2
View File
@@ -4,13 +4,14 @@ from datetime import datetime, timedelta, timezone
from uuid import UUID
from fastapi.testclient import TestClient
from sqlalchemy import create_engine
from sqlalchemy import create_engine, select
from sqlalchemy.orm import Session
from sqlalchemy.pool import StaticPool
from app.database import Base, get_session
from app.community_importer import stage_observations
from app.main import app
from app.models import Bait, BaitKind, CatchReport, Fish, ImportStatus, ModerationEvent, ModerationStatus, OfficialRecordImport, SourceType, Spot, Waterbody
from app.models import Bait, BaitKind, CatchReport, ExternalEntityAlias, ExternalObservation, Fish, ImportStatus, ModerationEvent, ModerationStatus, OfficialRecordImport, SourceType, Spot, Waterbody
engine = create_engine("sqlite://", connect_args={"check_same_thread": False}, poolclass=StaticPool)
@@ -89,6 +90,66 @@ def test_admin_requires_token() -> None:
assert client.get("/api/v1/admin/catch-reports").status_code == 401
assert client.get("/api/v1/admin/imports").status_code == 401
assert client.post("/api/v1/admin/imports/official-records").status_code == 401
assert client.get("/api/v1/admin/external-observations").status_code == 401
def test_external_observation_requires_mapping_and_complete_data_before_publication() -> None:
with Session(engine) as db:
stage_observations(db, [{
"source_system": "rf4db", "source_external_id": "review-complete",
"source_url": "https://rf4db.com/catches/review-complete", "fish": "Pike external",
"fish_external_id": "fish-1", "waterbody": "Lake external",
"waterbody_external_id": "lake-1", "x": 31, "y": 41, "weight_g": 6200,
"bait": "Тестовая приманка", "player_name": "External Player",
}])
observation_id = db.scalar(select(ExternalObservation.id).where(
ExternalObservation.source_external_id == "review-complete"
))
headers = {"Authorization": "Bearer change-me-in-production"}
premature = client.post(f"/api/v1/admin/external-observations/{observation_id}/publish", headers=headers)
assert premature.status_code == 409
mapped = client.patch(
f"/api/v1/admin/external-observations/{observation_id}/mapping", headers=headers,
json={"fish_slug": "pike", "waterbody_slug": "test-lake", "note": "verified fixture"},
)
assert mapped.status_code == 200
assert mapped.json()["status"] == "ready"
published = client.post(f"/api/v1/admin/external-observations/{observation_id}/publish", headers=headers)
assert published.status_code == 200
assert published.json()["status"] == "published"
repeated = client.post(f"/api/v1/admin/external-observations/{observation_id}/publish", headers=headers)
assert repeated.json()["catch_report_id"] == published.json()["catch_report_id"]
with Session(engine) as db:
observation = db.get(ExternalObservation, observation_id)
report = db.get(CatchReport, observation.catch_report_id)
aliases = list(db.scalars(select(ExternalEntityAlias).where(ExternalEntityAlias.source_system == "rf4db")))
assert report.moderation_status == ModerationStatus.approved
assert report.raw_payload["provenance"]["source_external_id"] == "review-complete"
assert {alias.entity_type for alias in aliases} == {"fish", "waterbody"}
def test_incomplete_external_observation_stays_out_of_public_data() -> None:
with Session(engine) as db:
stage_observations(db, [{
"source_system": "rf4db", "source_external_id": "review-incomplete",
"source_url": "https://rf4db.com/catches/review-incomplete", "fish": "Pike external",
"waterbody": "Lake external", "x": 32, "y": 42,
}])
observation_id = db.scalar(select(ExternalObservation.id).where(
ExternalObservation.source_external_id == "review-incomplete"
))
headers = {"Authorization": "Bearer change-me-in-production"}
mapped = client.patch(
f"/api/v1/admin/external-observations/{observation_id}/mapping", headers=headers,
json={"fish_slug": "pike", "waterbody_slug": "test-lake"},
)
assert mapped.json()["status"] == "mapped"
assert client.post(f"/api/v1/admin/external-observations/{observation_id}/publish", headers=headers).status_code == 409
rejected = client.patch(
f"/api/v1/admin/external-observations/{observation_id}/reject", headers=headers,
json={"reason": "weight is absent"},
)
assert rejected.json()["status"] == "rejected"
def test_admin_can_start_and_list_official_import(monkeypatch) -> None: