diff --git a/README.md b/README.md index f3c00c7..eefc450 100644 --- a/README.md +++ b/README.md @@ -18,6 +18,8 @@ RF4DB/RF4-STAT/RF4MAP/RF4 Posts сначала принимаются в изо Все пять community-парсеров подключены к отдельному scheduler-процессу. Состояние запусков и ошибок хранится в PostgreSQL, параллельный запуск одного источника блокируется, минимальный интервал жёстко ограничен 1800 секундами. Локально процесс включается профилем `docker compose --profile scheduler up -d`; detail-URL RF4MAP/RF4 Posts задаются переменными окружения. +После повторных ошибок scheduler увеличивает паузу экспоненциально до 24 часов и возвращается к 30 минутам после успеха. Публичная страница `/status` показывает свежесть и состояние источников без URL запросов, внутренних ошибок и другой диагностической информации. + Подробный план и актуальные чекбоксы находятся в [`docs/ROADMAP.md`](docs/ROADMAP.md). Результаты проверки интерфейса и пять приоритетных UX-пакетов описаны в [`docs/UI_UX_AUDIT.md`](docs/UI_UX_AUDIT.md). Production-контур для домена `rf4spotter.ru`, TLS, секреты, backup/restore и команды первого запуска описаны в [`deploy/README.md`](deploy/README.md). Он использует отдельный `compose.production.yaml`; локальный `compose.yaml` остаётся средой разработки. Production seed добавляет только справочники — демонстрационные уловы отключены. Изолированные проверки `deploy/test-production-bootstrap.sh` и `deploy/test-backup-restore.sh` подтверждают старт с пустых volumes и восстановление данных. diff --git a/apps/api/app/community_scheduler.py b/apps/api/app/community_scheduler.py index cd9dbdf..090d32b 100644 --- a/apps/api/app/community_scheduler.py +++ b/apps/api/app/community_scheduler.py @@ -16,6 +16,15 @@ from .logging_config import configure_logging from .models import CommunityImportRun, DataSource logger = logging.getLogger("rf4.community_scheduler") +MAX_BACKOFF_SECONDS = 24 * 60 * 60 + +def retry_delay(statuses: list[str]) -> int: + failures = 0 + for status in statuses: + if status != "failed": + break + failures += 1 + return min(settings.community_import_interval_seconds * (2 ** max(0, failures - 1)), MAX_BACKOFF_SECONDS) def configured_sources(): return { @@ -33,8 +42,10 @@ def run_source(source_system: str, *, now: datetime | None = None) -> bool: source = session.get(DataSource, source_system) if source is None or not source.enabled: return False - latest = session.scalar(select(CommunityImportRun.started_at).where(CommunityImportRun.source_system == source_system).order_by(CommunityImportRun.started_at.desc()).limit(1)) - if latest and (latest if latest.tzinfo else latest.replace(tzinfo=timezone.utc)) > current - timedelta(seconds=settings.community_import_interval_seconds): + 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): return False 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}"}): return False diff --git a/apps/api/app/main.py b/apps/api/app/main.py index 0da4a98..0033416 100644 --- a/apps/api/app/main.py +++ b/apps/api/app/main.py @@ -24,9 +24,9 @@ from .config import settings from .community_review import ExternalReviewError, map_observation, publish_observation, reject_observation from .importer import ImportAlreadyRunning, ImportSourceError, import_records, normalize from .logging_config import configure_logging -from .models import Bait, BaitKind, CatchReport, DataSource, ExternalObservation, Fish, ModerationEvent, ModerationStatus, OfficialRecordImport, SourceType, Spot, SubmissionAttempt, Waterbody +from .models import Bait, BaitKind, CatchReport, CommunityImportRun, DataSource, ExternalObservation, Fish, ModerationEvent, ModerationStatus, OfficialRecordImport, SourceType, Spot, SubmissionAttempt, Waterbody from .readiness import readiness_report -from .schemas import ActivityOut, AdminCatchReportOut, BaitOut, CatchOut, CatchReportAccepted, CatchReportCreate, CatchReportCreated, ExternalObservationDecision, ExternalObservationMapping, ExternalObservationOut, ExternalObservationPublished, FishOut, ImportRunOut, ModerationUpdate, OfficialRecordOut, PublicObservationOut, SpotOut, WaterbodyOut +from .schemas import ActivityOut, AdminCatchReportOut, BaitOut, CatchOut, CatchReportAccepted, CatchReportCreate, CatchReportCreated, ExternalObservationDecision, ExternalObservationMapping, ExternalObservationOut, ExternalObservationPublished, FishOut, ImportRunOut, ModerationUpdate, OfficialRecordOut, PublicObservationOut, SourceStatusOut, SpotOut, WaterbodyOut from .storage import ScreenshotError, client as storage_client, delete_screenshot, signed_screenshot_url, upload_screenshot @@ -209,6 +209,30 @@ def community_observations( return result +@app.get("/api/v1/source-status", response_model=list[SourceStatusOut]) +def source_status(db: Db) -> list[SourceStatusOut]: + now = datetime.now(timezone.utc) + result = [] + for source in db.scalars(select(DataSource).order_by(DataSource.name)): + runs = list(db.scalars(select(CommunityImportRun).where(CommunityImportRun.source_system == source.key).order_by(CommunityImportRun.started_at.desc()).limit(20))) + latest = runs[0] if runs else None + success = next((run for run in runs if run.status == "success"), None) + if not source.enabled: + state = "disabled" + elif latest is None: + state = "waiting" + elif latest.status == "failed": + state = "source_changed" if "CommunityParseError" in (latest.error_summary or "") else "temporarily_limited" + elif _aware(latest.started_at) < now - timedelta(seconds=settings.community_import_interval_seconds * 2): + state = "stale" + else: + state = "healthy" + result.append(SourceStatusOut(source_system=source.key, name=source.name, status=state, + last_started_at=latest.started_at if latest else None, last_success_at=success.started_at if success else None, + observations=db.scalar(select(func.count()).select_from(ExternalObservation).where(ExternalObservation.source_system == source.key)) or 0)) + return result + + @app.get("/api/v1/records", response_model=list[OfficialRecordOut]) def records( db: Db, fish: str | None = None, waterbody: str | None = None, diff --git a/apps/api/app/schemas.py b/apps/api/app/schemas.py index bc20fb9..33414c8 100644 --- a/apps/api/app/schemas.py +++ b/apps/api/app/schemas.py @@ -221,3 +221,12 @@ class ExternalObservationPublished(BaseModel): observation_id: UUID catch_report_id: UUID status: str + + +class SourceStatusOut(BaseModel): + source_system: str + name: str + status: str + last_started_at: datetime | None + last_success_at: datetime | None + observations: int diff --git a/apps/api/tests/test_api.py b/apps/api/tests/test_api.py index 65119e4..2d1a91a 100644 --- a/apps/api/tests/test_api.py +++ b/apps/api/tests/test_api.py @@ -12,7 +12,7 @@ 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 +from app.models import Bait, BaitKind, CatchReport, DataSource, ExternalEntityAlias, ExternalObservation, Fish, ImportStatus, ModerationEvent, ModerationStatus, OfficialRecordImport, SourceType, Spot, Waterbody engine = create_engine("sqlite://", connect_args={"check_same_thread": False}, poolclass=StaticPool) @@ -66,6 +66,17 @@ def test_list_pagination_and_filter_validation() -> None: assert client.get("/api/v1/admin/catch-reports?limit=101", headers=headers).status_code == 422 +def test_public_source_status_hides_internal_details() -> None: + with Session(engine) as db: + if db.get(DataSource, "rf4db") is None: + db.add(DataSource(key="rf4db", name="RF4DB", base_url="https://rf4db.com", default_confidence=70, enabled=True)) + db.commit() + response = client.get("/api/v1/source-status") + assert response.status_code == 200 + assert response.json() + assert all("error_summary" not in item and "source_url" not in item for item in response.json()) + + def test_liveness_does_not_probe_dependencies() -> None: response = client.get("/health?token=must-not-be-logged") assert response.json() == {"status": "ok"} diff --git a/apps/api/tests/test_community_scheduler.py b/apps/api/tests/test_community_scheduler.py index 49b8220..7326ad1 100644 --- a/apps/api/tests/test_community_scheduler.py +++ b/apps/api/tests/test_community_scheduler.py @@ -1,7 +1,7 @@ import pytest from pydantic import ValidationError -from app.community_scheduler import configured_sources +from app.community_scheduler import MAX_BACKOFF_SECONDS, configured_sources, retry_delay from app.config import Settings @@ -12,3 +12,10 @@ def test_all_authorized_sources_are_scheduled() -> None: def test_community_interval_cannot_be_less_than_30_minutes() -> None: with pytest.raises(ValidationError): Settings(community_import_interval_seconds=1799) + + +def test_failed_runs_back_off_but_success_resets_delay() -> None: + assert retry_delay(["failed"]) == 1800 + assert retry_delay(["failed", "failed", "failed"]) == 7200 + assert retry_delay(["failed"] * 20) == MAX_BACKOFF_SECONDS + assert retry_delay(["success", "failed"]) == 1800 diff --git a/apps/web/src/layouts/Layout.astro b/apps/web/src/layouts/Layout.astro index 6adc139..4e98588 100644 --- a/apps/web/src/layouts/Layout.astro +++ b/apps/web/src/layouts/Layout.astro @@ -1,6 +1,7 @@ --- import "../styles/global.css"; import "../styles/catalog.css"; +import "../styles/source-status.css"; import FishingIcon from "../components/FishingIcon.astro"; const { title = "RF4 Spotter — свежие точки клёва Russian Fishing 4", @@ -56,6 +57,6 @@ const jsonLd = JSON.stringify({
Свежие данные и честная оценка