Files
rf4-spotter/apps/api/tests/test_community_importer.py
T
ik 8b2d7e2e1c
CI / backend-and-migrations (push) Canceled after 0s
CI / astro-build (push) Canceled after 0s
CI / compose-e2e (push) Canceled after 0s
feat: suggest confirmed external aliases
2026-09-07 18:57:09 +07:00

205 lines
8.7 KiB
Python

import json
from dataclasses import asdict
from datetime import datetime, timezone
from pathlib import Path
import pytest
from sqlalchemy import create_engine, func, select
from sqlalchemy.orm import Session
from app.community_importer import CommunityImportError, stage_observations
from app.community_review import ExternalReviewError, map_observation, suggest_aliases
from app.database import Base
from app.models import CatchReport, DataSource, ExternalEntityAlias, ExternalObservation, Fish, Waterbody
from rf4_research.community_sources import parse_rf4db_catches, parse_rf4map_point, parse_rf4posts_spot
FIXTURE = Path(__file__).parents[3] / "tests" / "fixtures" / "rf4db_catches_sample.html"
FIXTURES = FIXTURE.parent
def record(source: str = "rf4db", external_id: str = "catch-1") -> dict[str, object]:
urls = {
"rf4db": f"https://rf4db.com/ru/catches/{external_id}",
"rf4stat-fishing": f"https://rf4-stat.ru/fishing/{external_id}",
"rf4map": f"https://rf4map.ru/points/{external_id}",
"rf4posts-spot": f"https://rf4-posts.com/ru/spots/{external_id}",
}
return {
"source_system": source,
"source_external_id": external_id,
"source_url": urls.get(source, f"https://rf4-stat.ru/fishing/{external_id}"),
"fish": "Щука",
"fish_external_id": "pike",
"waterbody": "Тестовое озеро",
"waterbody_external_id": "test-lake",
"x": 71,
"y": 92,
"weight_g": None,
"published_at": None,
"bait": "Приманка",
}
@pytest.fixture
def db() -> Session:
engine = create_engine("sqlite://")
Base.metadata.create_all(engine)
with Session(engine) as session:
yield session
engine.dispose()
def test_staging_is_idempotent_and_preserves_first_seen(db: Session) -> None:
first = datetime(2026, 9, 3, 10, tzinfo=timezone.utc)
second = datetime(2026, 9, 3, 11, tzinfo=timezone.utc)
assert stage_observations(db, [record()], fetched_at=first) == (1, 0)
changed = record() | {"x": 73, "weight_g": 5_000}
assert stage_observations(db, [changed], fetched_at=second) == (0, 1)
item = db.scalar(select(ExternalObservation))
assert item is not None
assert (item.x, item.weight_g, item.status) == (73, 5_000, "staged")
assert item.first_seen_at.replace(tzinfo=timezone.utc) == first
assert item.last_seen_at.replace(tzinfo=timezone.utc) == second
assert db.scalar(select(func.count()).select_from(ExternalObservation)) == 1
source = db.get(DataSource, "rf4db")
assert source is not None and source.enabled is True
def test_external_ids_are_isolated_by_source(db: Session) -> None:
created, updated = stage_observations(db, [record("rf4db"), record("rf4stat-fishing")])
assert (created, updated) == (2, 0)
def test_research_sources_enter_enabled_staging(db: Session) -> None:
created, updated = stage_observations(db, [record("rf4map"), record("rf4posts-spot")])
assert (created, updated) == (2, 0)
assert db.get(DataSource, "rf4map").enabled is True
assert db.get(DataSource, "rf4posts-spot").enabled is True
def test_complete_observation_with_reviewed_aliases_is_published(db: Session) -> None:
source = DataSource(key="rf4db", name="RF4DB", base_url="https://rf4db.com", default_confidence=70, enabled=True)
fish = Fish(slug="pike", name_ru="Щука")
waterbody = Waterbody(slug="test-lake", name_ru="Тестовое озеро")
db.add_all([source, fish, waterbody])
db.flush()
now = datetime.now(timezone.utc)
db.add_all([
ExternalEntityAlias(source_system="rf4db", entity_type="fish", external_id="pike", external_name="Щука", fish=fish, updated_at=now),
ExternalEntityAlias(source_system="rf4db", entity_type="waterbody", external_id="test-lake", external_name="Тестовое озеро", waterbody=waterbody, updated_at=now),
])
db.commit()
assert stage_observations(db, [record() | {"weight_g": 5_000}]) == (1, 0)
item = db.scalar(select(ExternalObservation))
assert item is not None
assert item.status == "published"
assert item.catch_report is not None
assert item.catch_report.fish_id == fish.id
assert item.catch_report.waterbody_id == waterbody.id
assert db.scalar(select(func.count()).select_from(CatchReport)) == 1
assert stage_observations(db, [record() | {"weight_g": 5_000}]) == (0, 1)
db.refresh(item)
assert item.status == "published"
assert db.scalar(select(func.count()).select_from(CatchReport)) == 1
def test_auto_publication_requires_enabled_source(db: Session) -> None:
source = DataSource(key="rf4db", name="RF4DB", base_url="https://rf4db.com", default_confidence=70, enabled=False)
fish = Fish(slug="pike", name_ru="Щука")
waterbody = Waterbody(slug="test-lake", name_ru="Тестовое озеро")
db.add_all([source, fish, waterbody])
db.flush()
now = datetime.now(timezone.utc)
db.add_all([
ExternalEntityAlias(source_system="rf4db", entity_type="fish", external_id="pike", external_name="Щука", fish=fish, updated_at=now),
ExternalEntityAlias(source_system="rf4db", entity_type="waterbody", external_id="test-lake", external_name="Тестовое озеро", waterbody=waterbody, updated_at=now),
])
db.commit()
stage_observations(db, [record() | {"weight_g": 5_000}])
item = db.scalar(select(ExternalObservation))
assert item is not None and item.status == "staged" and item.catch_report is None
def test_confirmed_aliases_are_suggestions_and_cannot_be_retargeted(db: Session) -> None:
source = DataSource(key="rf4db", name="RF4DB", base_url="https://rf4db.com", default_confidence=70, enabled=True)
pike = Fish(slug="pike", name_ru="Щука")
perch = Fish(slug="perch", name_ru="Окунь")
waterbody = Waterbody(slug="test-lake", name_ru="Тестовое озеро")
db.add_all([source, pike, perch, waterbody])
db.commit()
stage_observations(db, [record(external_id="first")])
first = db.scalar(select(ExternalObservation).where(ExternalObservation.source_external_id == "first"))
assert first is not None
map_observation(db, first, pike, waterbody)
stage_observations(db, [record(external_id="second")])
second = db.scalar(select(ExternalObservation).where(ExternalObservation.source_external_id == "second"))
assert second is not None
suggested_fish, suggested_waterbody = suggest_aliases(db, second)
assert (suggested_fish.slug, suggested_waterbody.slug) == ("pike", "test-lake")
with pytest.raises(ExternalReviewError, match="confirmed fish alias"):
map_observation(db, second, perch, waterbody)
db.rollback()
alias = db.scalar(select(ExternalEntityAlias).where(ExternalEntityAlias.entity_type == "fish"))
assert alias is not None and alias.fish_id == pike.id
def test_parser_json_can_be_staged_without_losing_provenance(db: Session) -> None:
parsed = parse_rf4db_catches(FIXTURE.read_text(encoding="utf-8"))
payload = json.loads(json.dumps([asdict(item) for item in parsed], default=str))
assert stage_observations(db, payload) == (1, 0)
item = db.scalar(select(ExternalObservation))
assert item is not None
assert (item.source_system, item.fish_external_id, item.x, item.y) == ("rf4db", "pike", 71, 92)
@pytest.mark.parametrize(("fixture_name", "source_url", "parser"), [
("rf4map_point_sample.html", "https://rf4map.ru/points/275", parse_rf4map_point),
(
"rf4posts_spot_sample.html",
"https://rf4-posts.com/ru/spots/d0c6d9c6-4ebf-49a7-98a8-9a562553a8ee",
parse_rf4posts_spot,
),
])
def test_new_parser_json_can_be_staged(
db: Session, fixture_name: str, source_url: str, parser,
) -> None:
parsed = parser((FIXTURES / fixture_name).read_text(encoding="utf-8"), source_url=source_url)
payload = json.loads(json.dumps([asdict(item) for item in parsed], default=str))
created, updated = stage_observations(db, payload)
assert (created, updated) == (len(parsed), 0)
assert db.scalar(select(func.count()).select_from(ExternalObservation)) == len(parsed)
def test_invalid_source_rolls_back_caller_transaction(db: Session) -> None:
with pytest.raises(CommunityImportError, match="unsupported source_system"):
stage_observations(db, [record("unknown")])
db.rollback()
assert db.scalar(select(func.count()).select_from(ExternalObservation)) == 0
@pytest.mark.parametrize("change", [
{"source_url": "https://attacker.example/catch-1"},
{"x": 10_001},
{"weight_g": 3_000_001},
])
def test_invalid_provenance_and_ranges_are_rejected(db: Session, change: dict[str, object]) -> None:
with pytest.raises(CommunityImportError):
stage_observations(db, [record() | change])
db.rollback()