feat: enable safe community auto publishing

This commit is contained in:
ik
2026-09-07 13:00:58 +07:00
parent c4d1c8b87b
commit 7217356594
10 changed files with 130 additions and 21 deletions
@@ -0,0 +1,23 @@
"""Enable all authorized community sources in the staging registry."""
from alembic import op
revision = "0012"
down_revision = "0011"
branch_labels = None
depends_on = None
def upgrade() -> None:
op.execute("""
INSERT INTO data_source (key, name, base_url, default_confidence, enabled) VALUES
('rf4db', 'RF4DB', 'https://rf4db.com', 70, true),
('rf4stat-fishing', 'RF4-STAT fishing', 'https://rf4-stat.ru/fishing/', 65, true),
('rf4stat-post', 'RF4-STAT posts', 'https://rf4-stat.ru/posts/', 60, true),
('rf4map', 'RF4MAP', 'https://rf4map.ru', 55, true),
('rf4posts-spot', 'RF4 Posts spots', 'https://rf4-posts.com', 50, true)
ON CONFLICT (key) DO UPDATE SET enabled = EXCLUDED.enabled
""")
def downgrade() -> None:
op.execute("UPDATE data_source SET enabled = false WHERE key IN ('rf4db', 'rf4stat-fishing', 'rf4stat-post', 'rf4map', 'rf4posts-spot')")
+39 -2
View File
@@ -8,7 +8,8 @@ from urllib.parse import urlparse
from sqlalchemy import select
from sqlalchemy.orm import Session
from .models import DataSource, ExternalObservation
from .community_review import publish_observation
from .models import DataSource, ExternalEntityAlias, ExternalObservation
SOURCE_DEFAULTS = {
@@ -36,6 +37,7 @@ def stage_observations(
) -> tuple[int, int]:
fetched_at = fetched_at or datetime.now(timezone.utc)
created = updated = 0
touched: list[ExternalObservation] = []
for raw in records:
payload = _json_payload(raw)
source_system = _required(payload, "source_system", 50)
@@ -45,7 +47,7 @@ def stage_observations(
source = session.get(DataSource, source_system)
if source is None:
name, base_url, confidence = SOURCE_DEFAULTS[source_system]
source = DataSource(key=source_system, name=name, base_url=base_url, default_confidence=confidence, enabled=False)
source = DataSource(key=source_system, name=name, base_url=base_url, default_confidence=confidence, enabled=True)
session.add(source)
session.flush()
observation = session.scalar(select(ExternalObservation).where(
@@ -78,10 +80,45 @@ def stage_observations(
for key, value in values.items():
setattr(observation, key, value)
updated += 1
touched.append(observation)
session.commit()
for observation in touched:
_auto_publish(session, observation)
return created, updated
def _auto_publish(session: Session, observation: ExternalObservation) -> bool:
"""Publish only complete observations covered by previously reviewed aliases."""
if (
observation.status not in {"staged", "mapped", "ready"}
or not observation.source.enabled
or observation.fish_external_id is None
or observation.waterbody_external_id is None
or observation.x is None
or observation.y is None
or observation.weight_g is None
):
return False
fish_alias = session.scalar(select(ExternalEntityAlias).where(
ExternalEntityAlias.source_system == observation.source_system,
ExternalEntityAlias.entity_type == "fish",
ExternalEntityAlias.external_id == observation.fish_external_id,
))
waterbody_alias = session.scalar(select(ExternalEntityAlias).where(
ExternalEntityAlias.source_system == observation.source_system,
ExternalEntityAlias.entity_type == "waterbody",
ExternalEntityAlias.external_id == observation.waterbody_external_id,
))
if fish_alias is None or fish_alias.fish is None or waterbody_alias is None or waterbody_alias.waterbody is None:
return False
observation.fish = fish_alias.fish
observation.waterbody = waterbody_alias.waterbody
observation.status = "ready"
observation.review_note = "Automatically matched by previously reviewed source aliases"
publish_observation(session, observation)
return True
def _json_payload(raw: dict[str, Any]) -> dict[str, Any]:
if not isinstance(raw, dict):
raise CommunityImportError("each observation must be an object")
+1 -1
View File
@@ -145,7 +145,7 @@ class DataSource(Base):
name: Mapped[str] = mapped_column(String(100))
base_url: Mapped[str] = mapped_column(Text)
default_confidence: Mapped[int]
enabled: Mapped[bool] = mapped_column(default=False)
enabled: Mapped[bool] = mapped_column(default=True)
class ExternalObservation(Base):
+53 -5
View File
@@ -9,7 +9,7 @@ from sqlalchemy.orm import Session
from app.community_importer import CommunityImportError, stage_observations
from app.database import Base
from app.models import DataSource, ExternalObservation
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
@@ -64,7 +64,7 @@ def test_staging_is_idempotent_and_preserves_first_seen(db: Session) -> None:
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 False
assert source is not None and source.enabled is True
def test_external_ids_are_isolated_by_source(db: Session) -> None:
@@ -73,12 +73,60 @@ def test_external_ids_are_isolated_by_source(db: Session) -> None:
assert (created, updated) == (2, 0)
def test_research_sources_can_enter_disabled_staging(db: Session) -> None:
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 False
assert db.get(DataSource, "rf4posts-spot").enabled is False
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_parser_json_can_be_staged_without_losing_provenance(db: Session) -> None: