diff --git a/README.md b/README.md index 85e8c15..1ee75f7 100644 --- a/README.md +++ b/README.md @@ -144,7 +144,7 @@ npm run test:e2e docker compose --profile tools run --rm importer ``` -Импорт делает до трёх ограниченных попыток, проверяет DOM-контракт и не удаляет ранее сохранённые данные при сбое. Повторный запуск обновляет совпавшие записи по SHA-256 ключу и не создаёт дубликаты. Расписание реализовано, но намеренно не включается обычным запуском: сначала требуется согласовать допустимость регулярного опроса официального сайта. +Импорт делает до трёх ограниченных попыток, проверяет DOM-контракт и не удаляет ранее сохранённые данные при сбое. Повторный запуск обновляет совпавшие записи по SHA-256 ключу и не создаёт дубликаты. PostgreSQL advisory lock не допускает параллельный импорт одной source/region/category через admin и scheduler. Расписание реализовано, но намеренно не включается обычным запуском: сначала требуется согласовать допустимость регулярного опроса официального сайта. Ручной административный запуск также доступен через `POST /api/v1/admin/imports/official-records`, журнал — через `GET /api/v1/admin/imports`. Импорт сохраняет HTTP-метаданные и использует `ETag`/`Last-Modified`, когда источник их предоставляет. diff --git a/apps/api/app/importer.py b/apps/api/app/importer.py index f2646da..ba7b887 100644 --- a/apps/api/app/importer.py +++ b/apps/api/app/importer.py @@ -3,11 +3,12 @@ from __future__ import annotations import hashlib import re import time as time_module +from contextlib import contextmanager from dataclasses import asdict, dataclass from datetime import date, datetime, time, timezone import httpx -from sqlalchemy import select +from sqlalchemy import select, text from sqlalchemy.orm import Session from rf4_research.official_parser import RecordsContractError, parse_official_records @@ -24,6 +25,10 @@ class ImportSourceError(ValueError): pass +class ImportAlreadyRunning(RuntimeError): + pass + + @dataclass(frozen=True, slots=True) class RawRecord: region: str @@ -105,7 +110,42 @@ def fetch_records( raise AssertionError("unreachable") +def _lock_key(url: str, region: str, category: str) -> int: + digest = hashlib.sha256(f"{url}|{region.upper()}|{category}".encode()).digest() + return int.from_bytes(digest[:8], byteorder="big", signed=True) + + +@contextmanager +def _official_import_lock(session: Session, *, url: str, region: str, category: str): + bind = session.get_bind() + if bind.dialect.name != "postgresql": + yield + return + connection = bind.connect() + key = _lock_key(url, region, category) + try: + acquired = bool(connection.scalar(text("SELECT pg_try_advisory_lock(:key)"), {"key": key})) + except Exception: + connection.close() + raise + if not acquired: + connection.close() + raise ImportAlreadyRunning("official import is already running for this source and category") + try: + yield + finally: + try: + connection.execute(text("SELECT pg_advisory_unlock(:key)"), {"key": key}) + finally: + connection.close() + + def import_records(session: Session, *, url: str, region: str, category: str, html: str | None = None) -> OfficialRecordImport: + with _official_import_lock(session, url=url, region=region, category=category): + return _import_records_locked(session, url=url, region=region, category=category, html=html) + + +def _import_records_locked(session: Session, *, url: str, region: str, category: str, html: str | None = None) -> OfficialRecordImport: run = OfficialRecordImport(started_at=datetime.now(timezone.utc), status=ImportStatus.running, source_url=url, rows_seen=0, rows_created=0, rows_updated=0) session.add(run) session.commit() diff --git a/apps/api/app/main.py b/apps/api/app/main.py index bb3b8b5..9e84f83 100644 --- a/apps/api/app/main.py +++ b/apps/api/app/main.py @@ -22,7 +22,7 @@ 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 .importer import ImportAlreadyRunning, ImportSourceError, import_records, normalize from .logging_config import configure_logging from .models import Bait, BaitKind, CatchReport, ExternalObservation, Fish, ModerationEvent, ModerationStatus, OfficialRecordImport, SourceType, Spot, SubmissionAttempt, Waterbody from .readiness import readiness_report @@ -200,6 +200,8 @@ def admin_start_official_import(db: Db, _: Annotated[str, Depends(_admin)]) -> O 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 diff --git a/apps/api/app/scheduler.py b/apps/api/app/scheduler.py index 6e7c582..4afaad3 100644 --- a/apps/api/app/scheduler.py +++ b/apps/api/app/scheduler.py @@ -9,7 +9,7 @@ from sqlalchemy.orm import Session from .config import settings from .database import SessionLocal -from .importer import import_records +from .importer import ImportAlreadyRunning, import_records from .logging_config import configure_logging from .models import OfficialRecordImport @@ -36,12 +36,16 @@ def run_due_import() -> bool: with SessionLocal() as session: if not import_is_due(session): return False - run = import_records( - session, - url=settings.official_records_url, - region=settings.official_records_region, - category=settings.official_records_category, - ) + try: + run = import_records( + session, + url=settings.official_records_url, + region=settings.official_records_region, + category=settings.official_records_category, + ) + except ImportAlreadyRunning: + logger.info("official import skipped because it is already running", extra={"event": "official_import_locked"}) + return False logger.info( "official import completed", extra={ diff --git a/apps/api/tests/test_api.py b/apps/api/tests/test_api.py index d5de88a..a2fda49 100644 --- a/apps/api/tests/test_api.py +++ b/apps/api/tests/test_api.py @@ -10,6 +10,7 @@ from sqlalchemy.pool import StaticPool from app.database import Base, get_session from app.community_importer import stage_observations +from app.importer import ImportAlreadyRunning from app.main import app from app.models import Bait, BaitKind, CatchReport, ExternalEntityAlias, ExternalObservation, Fish, ImportStatus, ModerationEvent, ModerationStatus, OfficialRecordImport, SourceType, Spot, Waterbody @@ -186,6 +187,12 @@ def test_admin_can_start_and_list_official_import(monkeypatch) -> None: listed = client.get("/api/v1/admin/imports?limit=1&offset=0", headers=headers) assert listed.status_code == 200 assert listed.json()[0]["id"] == started.json()["id"] + def busy_import(*args, **kwargs): + raise ImportAlreadyRunning("official import is already running") + + monkeypatch.setattr("app.main.import_records", busy_import) + conflict = client.post("/api/v1/admin/imports/official-records", headers=headers) + assert conflict.status_code == 409 def test_pending_report_accepts_one_validated_screenshot(monkeypatch) -> None: diff --git a/apps/api/tests/test_import_lock_postgres.py b/apps/api/tests/test_import_lock_postgres.py new file mode 100644 index 0000000..a8dc680 --- /dev/null +++ b/apps/api/tests/test_import_lock_postgres.py @@ -0,0 +1,27 @@ +import os + +import pytest +from sqlalchemy import create_engine +from sqlalchemy.orm import Session + +from app.importer import ImportAlreadyRunning, _official_import_lock + + +@pytest.mark.skipif(not os.environ.get("DATABASE_URL", "").startswith("postgresql"), reason="requires PostgreSQL") +def test_postgresql_import_lock_blocks_only_same_source_category() -> None: + engine = create_engine(os.environ["DATABASE_URL"]) + first = Session(engine) + second = Session(engine) + try: + with _official_import_lock(first, url="https://example.test/records", region="RU", category="records"): + with pytest.raises(ImportAlreadyRunning): + with _official_import_lock(second, url="https://example.test/records", region="RU", category="records"): + pass + with _official_import_lock(second, url="https://example.test/records", region="RU", category="weekly"): + pass + with _official_import_lock(second, url="https://example.test/records", region="RU", category="records"): + pass + finally: + first.close() + second.close() + engine.dispose() diff --git a/apps/api/tests/test_importer.py b/apps/api/tests/test_importer.py index 34977be..27b8eff 100644 --- a/apps/api/tests/test_importer.py +++ b/apps/api/tests/test_importer.py @@ -7,7 +7,7 @@ from sqlalchemy import create_engine, func, select from sqlalchemy.orm import Session from app.database import Base -from app.importer import FetchResult, ImportSourceError, import_records, parse_html +from app.importer import FetchResult, ImportAlreadyRunning, ImportSourceError, _lock_key, _official_import_lock, import_records, parse_html from app.models import CatchReport, ImportStatus, OfficialRecordImport, SourceType @@ -15,6 +15,26 @@ FIXTURE = Path(__file__).parents[3] / "tests" / "fixtures" / "records_ru_sample. WEEKLY_FIXTURE = Path(__file__).parents[3] / "tests" / "fixtures" / "weekly_records_sample.html" +def test_import_lock_is_stable_and_fails_closed_when_busy() -> None: + class Connection: + def scalar(self, statement, parameters): + assert "pg_try_advisory_lock" in str(statement) + assert parameters == {"key": _lock_key("https://example.test", "RU", "records")} + return False + + def close(self): + self.closed = True + + connection = Connection() + bind = type("Bind", (), {"dialect": type("Dialect", (), {"name": "postgresql"})(), "connect": lambda self: connection})() + session = type("Session", (), {"get_bind": lambda self: bind})() + assert _lock_key("https://example.test", "ru", "records") == _lock_key("https://example.test", "RU", "records") + with pytest.raises(ImportAlreadyRunning, match="already running"): + with _official_import_lock(session, url="https://example.test", region="RU", category="records"): + raise AssertionError("busy lock must not enter import") + assert connection.closed is True + + def test_parser_and_import_are_idempotent() -> None: html = FIXTURE.read_text(encoding="utf-8") parsed = parse_html(html, region="RU", category="records") diff --git a/deploy/README.md b/deploy/README.md index bf29083..7cbc893 100644 --- a/deploy/README.md +++ b/deploy/README.md @@ -63,6 +63,8 @@ docker compose --env-file .env.production -f compose.production.yaml exec api py Автоматический scheduler не входит в production-файл. RF4MAP и RF4 Posts нельзя опрашивать чаще одного раза в 30 минут; до отдельной эксплуатационной задачи используйте только контролируемые ручные запуски и staging. +Официальный импорт защищён PostgreSQL advisory lock на комбинацию source/region/category. Параллельный admin-запрос получает `409`, а scheduler записывает безопасный skip и не делает второй HTTP-запрос к источнику. + ## 5. Обновление ```bash diff --git a/docs/ROADMAP.md b/docs/ROADMAP.md index f06ee7f..0c8f8fa 100644 --- a/docs/ROADMAP.md +++ b/docs/ROADMAP.md @@ -2,7 +2,7 @@ Этот файл — рабочий источник правды по развитию проекта. После завершения задачи её чекбокс меняется с `[ ]` на `[x]`, рядом добавляется ссылка на коммит или короткое подтверждение проверки. Новые задачи добавляются в соответствующий этап, а не хранятся только в переписке. -Последняя сверка плана со спецификацией, кодом и UI/UX-аудитом: 6 сентября 2026 года. +Последняя сверка плана со спецификацией, кодом и UI/UX-аудитом: 7 сентября 2026 года. Обозначения: @@ -24,7 +24,7 @@ - [x] Привести журнал импорта к административному контракту `GET /api/v1/admin/imports` с авторизацией, пагинацией и стабильной сортировкой (проверено API-тестом). - [x] Добавить HTTP-кэширование источника (`ETag`/`Last-Modified`, если источник их отдаёт) и сохранить диагностические метаданные ответа (миграция `0005`, тест условного запроса и `304`). - [x] Добавить планировщик импорта с безопасной частотой по умолчанию один раз в 60 минут; отдельный opt-in контейнер/процесс (профиль `scheduler`, обычным запуском не активируется). -- [ ] Защитить официальный импорт PostgreSQL advisory lock или эквивалентом, чтобы ручной endpoint и несколько scheduler-процессов не импортировали одну категорию одновременно. +- [x] Защитить официальный импорт PostgreSQL session advisory lock: одинаковая source/region/category не запускается параллельно, admin получает `409`, scheduler безопасно пропускает цикл; проверено двумя независимыми PostgreSQL-соединениями. - [x] Проверить актуальные `robots.txt` и условия использования перед включением расписания; результат записать в `docs/data-sources.md` (`robots.txt` вернул `404`; автоматический профиль оставлен выключенным до явного разрешения). - [x] Добавить интеграционные тесты: повторный импорт не создаёт дубликаты, сбой источника не удаляет данные, изменение DOM завершается понятной ошибкой.