diff --git a/README.md b/README.md index 2f13fe1..203e7f2 100644 --- a/README.md +++ b/README.md @@ -22,7 +22,7 @@ Web Docker-образ устанавливает зависимости чере Функциональный MVP и локальный production-контур готовятся к открытой альфе: официальный импорт, пользовательские заявки, модерация, объяснимый индекс, staging внешних источников, адаптивный Astro UI, миграции, резервное копирование, retention, мониторинг и security/accessibility-проверки реализованы. На всех страницах подключён компактный баннер открытой альфы со ссылками на статус, правила и отправку улова. В production Compose включён community scheduler; локально он запускается отдельным профилем. Публичный запуск блокируют покупка и настройка сервера, DNS/TLS, реальные секреты, внешний backup, канал уведомлений; публичный адрес обратной связи ещё не задан. -RF4DB/RF4-STAT/RF4MAP/RF4 Posts сначала принимаются в изолированный staging. Полные записи с ранее подтверждёнными алиасами источника публикуются автоматически; новые соответствия и неполные записи остаются на ручной проверке. Admin API предлагает точные ранее подтверждённые алиасы отдельно от mapping-действия и запрещает молча переназначать alias другой сущности. Для разрешённых community-источников действует интервал не менее 30 минут на источник. Открытая альфа не использует продуктовый allowlist: интерфейс показывает весь корректно загруженный разрешённый каталог, сохраняя требования полноты и модерации. +RF4DB/RF4-STAT/RF4MAP/RF4 Posts сначала принимаются в изолированный staging. Полные записи с ранее подтверждёнными алиасами источника публикуются автоматически; новые соответствия и неполные записи остаются на ручной проверке. Admin API предлагает точные ранее подтверждённые алиасы отдельно от mapping-действия и запрещает молча переназначать alias другой сущности. Для разрешённых community-источников действует интервал не менее 30 минут на сайт, общий для всех его endpoint. Открытая альфа не использует продуктовый allowlist: интерфейс показывает весь корректно загруженный разрешённый каталог, сохраняя требования полноты и модерации. На сайте у каждой записи отображается источник, а у агрегированной активности — все вошедшие в расчёт источники. Неполные community-наблюдения публикуются сразу в отдельной ленте «Полевые сигналы» с предупреждением и перечнем отсутствующих полей; до подтверждения полноты они не влияют на индекс клёва. Лента раскрывается серверной кнопкой «Показать ещё», сохраняет выбранные фильтры и ограничена 48 сигналами на страницу. Визуально объединяются только повторы одного ID источника; похожие записи разных площадок остаются самостоятельными наблюдениями. Sidebar лидера скрывается при единственном результате, чтобы не повторять ту же карточку. @@ -30,7 +30,7 @@ RF4DB/RF4-STAT/RF4MAP/RF4 Posts сначала принимаются в изо Публичные точки используют постоянные читаемые адреса вида `/spots/kuori-85x92`; старые UUID-адреса остаются совместимыми и перенаправляются на канонический URL. На странице точки координаты дополнительно показаны фирменным радаром, который не имитирует отсутствующую географию водоёма, а уловы за 72 часа — шкалой-леской с 12-часовым шагом. Каждый улов показывает источник, относительную свежесть и точное время UTC; время получения явно отделено от времени улова. Карточки активности и каталог дополнены лёгкими SVG-силуэтами рыб без внешних графических зависимостей. Пустые и аварийные состояния используют собственную CSS-иллюстрацию поплавка; анимация учитывает системное ограничение движения. -Все пять community-парсеров подключены к отдельному scheduler-процессу. Состояние запусков и ошибок хранится в PostgreSQL, параллельный запуск одного источника блокируется, минимальный интервал жёстко ограничен 1800 секундами. Локально процесс включается профилем `docker compose --profile scheduler up -d`; detail-URL RF4MAP/RF4 Posts задаются переменными окружения. +Все пять community-парсеров подключены к отдельному scheduler-процессу. Попытка резервируется в PostgreSQL до HTTP-запроса, поэтому ошибки тоже расходуют cooldown. Блокировка и минимальный интервал 1800 секунд действуют на весь домен; endpoint одного сайта выбираются по самому давнему запуску и не голодают. Ручной production-запуск использует тот же журнал: `docker compose exec api python -m app.cli fetch-community rf4stat-fishing`. Локально scheduler включается профилем `docker compose --profile scheduler up -d`; detail-URL RF4MAP/RF4 Posts задаются переменными окружения. После повторных ошибок scheduler увеличивает паузу экспоненциально до 24 часов и возвращается к 30 минутам после успеха. Публичная страница `/status` показывает свежесть и состояние источников без URL запросов, внутренних ошибок и другой диагностической информации. @@ -73,7 +73,7 @@ python -m rf4_research.community_cli rf4map-point --url https://rf4map.ru/points python -m rf4_research.community_cli rf4posts-spot --url https://rf4-posts.com/ru/spots/UUID --limit 25 ``` -Команды печатают нормализованный JSON в stdout и ничего не записывают в базу. Detail-команды требуют явный публичный URL и не обходят запрещённые `/api/`. Для RF4-STAT действует пауза не менее пяти секунд между разными страницами; для RF4MAP/RF4 Posts CLI хранит состояние в `.cache/community-fetch-state.json` и блокирует повтор того же источника раньше 30 минут. +Команды печатают нормализованный JSON в stdout и ничего не записывают в базу. Detail-команды требуют явный публичный URL и не обходят запрещённые `/api/`. CLI резервирует домен в `.cache/community-fetch-state.json` до HTTP-запроса и блокирует любой его endpoint на 30 минут даже после ошибки. Это автономный исследовательский режим: не запускайте его одновременно с production scheduler; для ручного production-запуска используйте `app.cli fetch-community`, который разделяет PostgreSQL-cooldown с scheduler. Проверенный JSON можно идемпотентно загрузить в изолированный staging, не влияющий на публичную статистику: diff --git a/apps/api/app/cli.py b/apps/api/app/cli.py index f24138c..ee68d5c 100644 --- a/apps/api/app/cli.py +++ b/apps/api/app/cli.py @@ -12,6 +12,7 @@ from .community_importer import stage_observations from .retention import RetentionPolicy, apply_retention from .storage import delete_screenshot from .catalog_audit import audit_catalog +from .community_scheduler import configured_sources, run_source def main() -> int: @@ -24,6 +25,8 @@ def main() -> int: community = sub.add_parser("stage-community-json") community.add_argument("--input", default="-", help="JSON array path or - for stdin") community.add_argument("--limit", type=int, default=500) + fetch_community = sub.add_parser("fetch-community") + fetch_community.add_argument("source", choices=configured_sources()) cleanup = sub.add_parser("cleanup-retention") cleanup.add_argument("--apply", action="store_true", help="apply changes; default is dry-run") sub.add_parser("audit-catalog") @@ -45,6 +48,9 @@ def main() -> int: parser.error("input must be a JSON array") created, updated = stage_observations(session, payload[:args.limit]) print(f"staged: created={created} updated={updated}") + elif args.command == "fetch-community": + started = run_source(args.source) + print("community fetch started" if started else "community fetch skipped: disabled, locked, or cooling down") elif args.command == "cleanup-retention": policy = RetentionPolicy( submission_days=settings.retention_submission_days, diff --git a/apps/api/app/community_scheduler.py b/apps/api/app/community_scheduler.py index 98d04ad..76ff216 100644 --- a/apps/api/app/community_scheduler.py +++ b/apps/api/app/community_scheduler.py @@ -7,7 +7,7 @@ from datetime import datetime, timedelta, timezone from sqlalchemy import select, text -from rf4_research.community_cli import SOURCES, fetch_html +from rf4_research.community_cli import SOURCES, fetch_html, fetch_site_key from rf4_research.community_sources import parse_rf4map_point, parse_rf4posts_spot from .community_importer import stage_observations from .config import settings @@ -35,18 +35,32 @@ def configured_sources(): "rf4posts-spot": (settings.rf4posts_spot_url, parse_rf4posts_spot), } +def oldest_site_source(source_system: str, latest_by_source: dict[str, datetime]) -> str: + sources = configured_sources() + site = fetch_site_key(sources[source_system][0]) + candidates = [key for key, (url, _) in sources.items() if fetch_site_key(url) == site] + order = {key: index for index, key in enumerate(candidates)} + return min(candidates, key=lambda key: (latest_by_source.get(key, datetime.min.replace(tzinfo=timezone.utc)), order[key])) + def run_source(source_system: str, *, now: datetime | None = None) -> bool: current = now or datetime.now(timezone.utc) url, parser = configured_sources()[source_system] + site_key = fetch_site_key(url) + site_sources = [key for key, (candidate_url, _) in configured_sources().items() if fetch_site_key(candidate_url) == site_key] with SessionLocal() as session: source = session.get(DataSource, source_system) if source is None or not source.enabled: return False # Lock before reading cooldown: committing the reservation makes it visible # to the next contender before releasing this transaction lock. - if session.bind and session.bind.dialect.name == "postgresql" and not session.scalar(text("select pg_try_advisory_xact_lock(hashtext(:key))"), {"key": f"community:{source_system}"}): + if session.bind and session.bind.dialect.name == "postgresql" and not session.scalar(text("select pg_try_advisory_xact_lock(hashtext(:key))"), {"key": f"community-site:{site_key}"}): + return False + recent = list(session.scalars(select(CommunityImportRun).where(CommunityImportRun.source_system.in_(site_sources)).order_by(CommunityImportRun.started_at.desc()).limit(32))) + latest_by_source: dict[str, datetime] = {} + for previous in recent: + latest_by_source.setdefault(previous.source_system, previous.started_at if previous.started_at.tzinfo else previous.started_at.replace(tzinfo=timezone.utc)) + if oldest_site_source(source_system, latest_by_source) != source_system: return False - recent = list(session.scalars(select(CommunityImportRun).where(CommunityImportRun.source_system == source_system).order_by(CommunityImportRun.started_at.desc()).limit(8))) latest = recent[0].started_at if recent else None delay = retry_delay([run.status for run in recent]) if latest and (latest if latest.tzinfo else latest.replace(tzinfo=timezone.utc)) > current - timedelta(seconds=delay): diff --git a/apps/api/tests/test_community_scheduler.py b/apps/api/tests/test_community_scheduler.py index 7326ad1..67b27da 100644 --- a/apps/api/tests/test_community_scheduler.py +++ b/apps/api/tests/test_community_scheduler.py @@ -1,7 +1,9 @@ import pytest from pydantic import ValidationError -from app.community_scheduler import MAX_BACKOFF_SECONDS, configured_sources, retry_delay +from datetime import datetime, timedelta, timezone + +from app.community_scheduler import MAX_BACKOFF_SECONDS, configured_sources, oldest_site_source, retry_delay from app.config import Settings @@ -19,3 +21,10 @@ def test_failed_runs_back_off_but_success_resets_delay() -> None: assert retry_delay(["failed", "failed", "failed"]) == 7200 assert retry_delay(["failed"] * 20) == MAX_BACKOFF_SECONDS assert retry_delay(["success", "failed"]) == 1800 + + +def test_same_site_endpoints_rotate_by_oldest_attempt() -> None: + now = datetime.now(timezone.utc) + assert oldest_site_source("rf4stat-fishing", {}) == "rf4stat-fishing" + latest = {"rf4stat-fishing": now, "rf4stat-post": now - timedelta(hours=1)} + assert oldest_site_source("rf4stat-fishing", latest) == "rf4stat-post" diff --git a/docs/AUDIT_FIXES.md b/docs/AUDIT_FIXES.md index 06f5efd..e7141e1 100644 --- a/docs/AUDIT_FIXES.md +++ b/docs/AUDIT_FIXES.md @@ -19,7 +19,7 @@ ## Остаётся -- [ ] Единый лимит по сайту для CLI и scheduler, включая неуспешные попытки. +- [x] Production CLI и scheduler используют общий PostgreSQL-журнал и блокировку по домену; попытка резервируется до HTTP, включая ошибки. Endpoint одного сайта ротируются по давности. Автономный research CLI также ограничивает домен и ошибки, но явно запрещён параллельно с production scheduler (8 сентября 2026). - [ ] Межпроцессная инвалидация: сейчас фоновая публикация видна после TTL. - [x] Изменённая публикация снимается с активности и отправляется на ручное сопоставление. Подтверждение обновляет прежний CatchReport без дубликата; старый снимок сохраняется до подтверждения. Изменения видны в API после TTL кэша. - [ ] Обнаружение удалённых оригиналов и долговременная история всех редакций источника. diff --git a/docs/data-source-audit.md b/docs/data-source-audit.md index f16f4f8..6fa559b 100644 --- a/docs/data-source-audit.md +++ b/docs/data-source-audit.md @@ -65,7 +65,7 @@ Telegram, Discord и VK могут давать свежие координат Публичная detail-страница RF4 Posts содержит устойчивый UUID точки, координаты, slug водоёма, список slug рыб, способ ловли, оснастку, клипсу, дату и ссылки на доказательства. Русская локализация позволяет связать slug с отображаемым названием. Контрольный пост дал **6 записей видов рыб из одной точки**. Это инструкция по точке, а не шесть доказанных индивидуальных уловов, поэтому вес остаётся `null`, а происхождение сохраняет общий UUID поста. -Оба fail-closed парсера добавлены в read-only исследовательский CLI и разрешены владельцем проекта с интервалом не менее 30 минут на источник. Источники зарегистрированы для изолированного staging выключенными по умолчанию. CLI хранит время последнего успешного получения и отклоняет слишком ранний повтор. Автоматического расписания и публикации нет; RF4 Posts считается агрегированной точкой, а не набором взвешенных уловов. +Оба fail-closed парсера добавлены в read-only исследовательский CLI и разрешены владельцем проекта с интервалом не менее 30 минут на сайт. CLI резервирует домен до запроса и отклоняет слишком ранний повтор любого endpoint, включая повтор после ошибки. В production все разрешённые адаптеры подключены к scheduler и staging; полные записи с подтверждёнными алиасами публикуются автоматически, остальные требуют проверки. RF4 Posts считается агрегированной точкой, а не набором взвешенных уловов. Также проверены два менее пригодных кандидата: diff --git a/rf4_research/community_cli.py b/rf4_research/community_cli.py index f7d3b1f..3c2ad5c 100644 --- a/rf4_research/community_cli.py +++ b/rf4_research/community_cli.py @@ -7,6 +7,7 @@ import sys import time from dataclasses import asdict from pathlib import Path +from urllib.parse import urlsplit from urllib.request import Request, urlopen from .community_sources import ( @@ -32,6 +33,16 @@ MIN_FETCH_INTERVAL_SECONDS = 30 * 60 DEFAULT_STATE_FILE = Path(".cache/community-fetch-state.json") +def fetch_site_key(url: str) -> str: + """Return a stable cooldown key shared by all endpoints of one site.""" + hostname = (urlsplit(url).hostname or "").lower() + if hostname.startswith("www."): + hostname = hostname[4:] + if not hostname: + raise ValueError("source URL must include a hostname") + return hostname + + def enforce_fetch_interval( source: str, *, state_file: Path, now: float | None = None, ) -> None: @@ -83,9 +94,11 @@ def main(argv: list[str] | None = None) -> int: default_url, parse = SOURCES.get(args.source, (None, DETAIL_SOURCES.get(args.source))) url = args.url or default_url try: - enforce_fetch_interval(args.source, state_file=args.state_file) + site_key = fetch_site_key(url) + enforce_fetch_interval(site_key, state_file=args.state_file) + # Reserve before network I/O: failed attempts count toward the limit too. + mark_fetch(site_key, state_file=args.state_file) html = fetch_html(url) - mark_fetch(args.source, state_file=args.state_file) records = (parse(html, source_url=url) if args.source in DETAIL_SOURCES else parse(html))[:args.limit] except Exception as exc: print(f"community source failed: {exc}", file=sys.stderr) diff --git a/tests/test_community_cli.py b/tests/test_community_cli.py index afcb9f0..e18d406 100644 --- a/tests/test_community_cli.py +++ b/tests/test_community_cli.py @@ -1,8 +1,10 @@ +import json from pathlib import Path import pytest -from rf4_research.community_cli import enforce_fetch_interval, mark_fetch +from rf4_research import community_cli +from rf4_research.community_cli import enforce_fetch_interval, fetch_site_key, mark_fetch def test_fetch_cooldown_is_persistent_per_source(tmp_path: Path) -> None: @@ -13,3 +15,19 @@ def test_fetch_cooldown_is_persistent_per_source(tmp_path: Path) -> None: enforce_fetch_interval("rf4map-point", state_file=state_file, now=1_000) enforce_fetch_interval("rf4posts-spot", state_file=state_file, now=1_000) enforce_fetch_interval("rf4map-point", state_file=state_file, now=2_800) + + +def test_fetch_site_key_groups_endpoints_and_normalizes_www() -> None: + assert fetch_site_key("https://rf4-stat.ru/fishing/") == "rf4-stat.ru" + assert fetch_site_key("https://www.rf4-stat.ru/posts/") == "rf4-stat.ru" + + +def test_failed_fetch_still_reserves_site_cooldown(tmp_path: Path, monkeypatch: pytest.MonkeyPatch) -> None: + state_file = tmp_path / "fetch-state.json" + + def fail(_url: str) -> str: + raise OSError("offline") + + monkeypatch.setattr(community_cli, "fetch_html", fail) + assert community_cli.main(["rf4db", "--state-file", str(state_file)]) == 1 + assert "download.rf4db.com" in json.loads(state_file.read_text(encoding="utf-8"))