diff --git a/apps/api/app/community_scheduler.py b/apps/api/app/community_scheduler.py index fb913c9..ea34482 100644 --- a/apps/api/app/community_scheduler.py +++ b/apps/api/app/community_scheduler.py @@ -26,23 +26,31 @@ def retry_delay(statuses: list[str]) -> int: failures += 1 return min(settings.community_import_interval_seconds * (2 ** max(0, failures - 1)), MAX_BACKOFF_SECONDS) -def configured_sources(): - with SessionLocal() as session: - enabled_keys = { - s.key for s in session.scalars(select(DataSource).where(DataSource.enabled.is_(True))) - } + +def _static_registry() -> dict[str, tuple[str, callable]]: + """Return the full static registry without DB access (for unit tests).""" return { - k: v for k, v in { - "rf4db": SOURCES["rf4db"], - "rf4stat-fishing": SOURCES["rf4stat-fishing"], - "rf4stat-post": (SOURCES["rf4stat-posts"][0], SOURCES["rf4stat-posts"][1]), - "rf4map": (settings.rf4map_point_url, parse_rf4map_point), - "rf4posts-spot": (settings.rf4posts_spot_url, parse_rf4posts_spot), - }.items() if k in enabled_keys + "rf4db": SOURCES["rf4db"], + "rf4stat-fishing": SOURCES["rf4stat-fishing"], + "rf4stat-post": (SOURCES["rf4stat-posts"][0], SOURCES["rf4stat-posts"][1]), + "rf4map": (settings.rf4map_point_url, parse_rf4map_point), + "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() + +def configured_sources(enabled_keys: set[str] | None = None) -> dict[str, tuple[str, callable]]: + """Return enabled sources. When enabled_keys is None, query the DB.""" + registry = _static_registry() + if enabled_keys is None: + with SessionLocal() as session: + enabled_keys = { + s.key for s in session.scalars(select(DataSource).where(DataSource.enabled.is_(True))) + } + return {k: v for k, v in registry.items() if k in enabled_keys} + + +def oldest_site_source(source_system: str, latest_by_source: dict[str, datetime], enabled_keys: set[str] | None = None) -> str: + sources = configured_sources(enabled_keys) 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)} @@ -50,9 +58,12 @@ def oldest_site_source(source_system: str, latest_by_source: dict[str, datetime] def run_source(source_system: str, *, now: datetime | None = None) -> bool: current = now or datetime.now(timezone.utc) - url, parser = configured_sources()[source_system] + with SessionLocal() as session: + enabled_keys = {s.key for s in session.scalars(select(DataSource).where(DataSource.enabled.is_(True)))} + registry = _static_registry() + url, parser = registry[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] + site_sources = [key for key, (candidate_url, _) in registry.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: @@ -65,7 +76,7 @@ def run_source(source_system: str, *, now: datetime | None = None) -> bool: 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: + if oldest_site_source(source_system, latest_by_source, enabled_keys) != source_system: return False latest = recent[0].started_at if recent else None delay = retry_delay([run.status for run in recent]) diff --git a/apps/api/tests/test_api.py b/apps/api/tests/test_api.py index f8abfc4..712f734 100644 --- a/apps/api/tests/test_api.py +++ b/apps/api/tests/test_api.py @@ -45,11 +45,15 @@ def test_activity_filters_and_explains_score() -> None: response = client.get("/api/v1/activity?waterbody=test-lake&fish=pike&hours=24") assert response.status_code == 200 payload = response.json() - assert len(payload) == 1 - assert payload[0]["catches"] == 3 - assert payload[0]["unique_players"] == 3 - assert "3 свежих улова" in payload[0]["explanation"] - assert payload[0]["sources"] == ["manual-import"] + assert "items" in payload + assert payload["total"] == 1 + assert payload["limit"] == 20 + assert payload["offset"] == 0 + assert len(payload["items"]) == 1 + assert payload["items"][0]["catches"] == 3 + assert payload["items"][0]["unique_players"] == 3 + assert "3 свежих улова" in payload["items"][0]["explanation"] + assert payload["items"][0]["sources"] == ["manual-import"] def test_invalid_period_is_rejected() -> None: @@ -157,7 +161,7 @@ def test_admin_diagnostics_exposes_build_identity_only_to_admin() -> None: def test_spot_detail_and_catches() -> None: - spot_id = client.get("/api/v1/activity").json()[0]["spot_id"] + spot_id = client.get("/api/v1/activity").json()["items"][0]["spot_id"] detail = client.get(f"/api/v1/spots/{spot_id}") catches = client.get(f"/api/v1/spots/{spot_id}/catches") assert detail.status_code == 200 @@ -206,7 +210,7 @@ def test_user_report_requires_moderation_before_activity() -> None: approved = client.patch(f"/api/v1/admin/catch-reports/{report_id}", headers=headers, json={"status": "approved", "reason": "fixture verified"}) assert approved.status_code == 200 activity = client.get("/api/v1/activity?waterbody=test-lake&fish=pike&hours=24").json() - assert any(item["x"] == 77 and item["catches"] == 1 for item in activity) + assert any(item["x"] == 77 and item["catches"] == 1 for item in activity["items"]) def test_admin_requires_token() -> None: @@ -267,7 +271,7 @@ def test_incomplete_external_observation_is_publicly_labelled_but_not_counted() assert signal["source_system"] == "rf4db" assert signal["quality"] == "incomplete" assert signal["missing_fields"] == ["вес"] - assert all(item["x"] != 32 or item["y"] != 42 for item in client.get("/api/v1/activity").json()) + assert all(item["x"] != 32 or item["y"] != 42 for item in client.get("/api/v1/activity").json()["items"]) headers = {"Authorization": "Bearer change-me-in-production"} mapped = client.patch( f"/api/v1/admin/external-observations/{observation_id}/mapping", headers=headers, diff --git a/apps/api/tests/test_community_scheduler.py b/apps/api/tests/test_community_scheduler.py index 67b27da..366e18c 100644 --- a/apps/api/tests/test_community_scheduler.py +++ b/apps/api/tests/test_community_scheduler.py @@ -3,12 +3,12 @@ from pydantic import ValidationError from datetime import datetime, timedelta, timezone -from app.community_scheduler import MAX_BACKOFF_SECONDS, configured_sources, oldest_site_source, retry_delay +from app.community_scheduler import MAX_BACKOFF_SECONDS, configured_sources, _static_registry, oldest_site_source, retry_delay from app.config import Settings def test_all_authorized_sources_are_scheduled() -> None: - assert set(configured_sources()) == {"rf4db", "rf4stat-fishing", "rf4stat-post", "rf4map", "rf4posts-spot"} + assert set(_static_registry()) == {"rf4db", "rf4stat-fishing", "rf4stat-post", "rf4map", "rf4posts-spot"} def test_community_interval_cannot_be_less_than_30_minutes() -> None: @@ -25,6 +25,7 @@ def test_failed_runs_back_off_but_success_resets_delay() -> None: def test_same_site_endpoints_rotate_by_oldest_attempt() -> None: now = datetime.now(timezone.utc) - assert oldest_site_source("rf4stat-fishing", {}) == "rf4stat-fishing" + all_keys = {"rf4db", "rf4stat-fishing", "rf4stat-post", "rf4map", "rf4posts-spot"} + assert oldest_site_source("rf4stat-fishing", {}, all_keys) == "rf4stat-fishing" latest = {"rf4stat-fishing": now, "rf4stat-post": now - timedelta(hours=1)} - assert oldest_site_source("rf4stat-fishing", latest) == "rf4stat-post" + assert oldest_site_source("rf4stat-fishing", latest, all_keys) == "rf4stat-post" diff --git a/apps/api/tests/test_rate_limit.py b/apps/api/tests/test_rate_limit.py index 4c4dabc..b0bf55c 100644 --- a/apps/api/tests/test_rate_limit.py +++ b/apps/api/tests/test_rate_limit.py @@ -1,6 +1,7 @@ from __future__ import annotations from datetime import datetime, timezone +from unittest.mock import MagicMock import pytest from fastapi import HTTPException @@ -17,11 +18,34 @@ def test_rate_limit_is_persistent_and_does_not_store_raw_client() -> None: Base.metadata.create_all(engine) with Session(engine) as db: for _ in range(5): - _check_rate_limit("203.0.113.42", db) + mock_request = MagicMock() + mock_request.client.host = "203.0.113.42" + mock_request.headers.get.return_value = None + _check_rate_limit(mock_request, db) with pytest.raises(HTTPException) as blocked: - _check_rate_limit("203.0.113.42", db) + mock_request = MagicMock() + mock_request.client.host = "203.0.113.42" + mock_request.headers.get.return_value = None + _check_rate_limit(mock_request, db) assert blocked.value.status_code == 429 attempts = list(db.scalars(select(SubmissionAttempt))) assert len(attempts) == 5 assert all(item.client_hash != "203.0.113.42" and len(item.client_hash) == 64 for item in attempts) assert all(item.created_at.replace(tzinfo=timezone.utc) <= datetime.now(timezone.utc) for item in attempts) + + +def test_rate_limit_uses_forwarded_for_header() -> None: + engine = create_engine("sqlite://") + Base.metadata.create_all(engine) + with Session(engine) as db: + mock_real = MagicMock() + mock_real.client.host = "10.0.0.1" + mock_real.headers.get.return_value = "198.51.100.10" + for _ in range(5): + _check_rate_limit(mock_real, db) + with pytest.raises(HTTPException) as blocked: + mock_other = MagicMock() + mock_other.client.host = "10.0.0.2" + mock_other.headers.get.return_value = "198.51.100.10" + _check_rate_limit(mock_other, db) + assert blocked.value.status_code == 429 diff --git a/apps/api/tests/test_readiness.py b/apps/api/tests/test_readiness.py index 02252a5..2187622 100644 --- a/apps/api/tests/test_readiness.py +++ b/apps/api/tests/test_readiness.py @@ -28,11 +28,11 @@ def test_optional_import_does_not_block_dependencies() -> None: session, AvailableStorage(), import_required=False, import_interval_seconds=3600, ) assert ready is True - assert components == { - "postgresql": {"status": "ready"}, - "minio": {"status": "ready"}, - "official_import": {"status": "optional", "last_run_status": None}, - } + assert components["postgresql"]["status"] == "ready" + assert components["minio"]["status"] == "ready" + assert components["official_import"]["status"] == "optional" + assert components["official_import"]["last_run_status"] is None + assert "community_scheduler" in components def test_required_import_must_be_recent_and_successful() -> None: diff --git a/apps/web/src/pages/fish/[slug].astro b/apps/web/src/pages/fish/[slug].astro index ba4cb3d..94fdbe6 100644 --- a/apps/web/src/pages/fish/[slug].astro +++ b/apps/web/src/pages/fish/[slug].astro @@ -1,13 +1,13 @@ --- import ActivityCard from "../../components/ActivityCard.astro"; import Layout from "../../layouts/Layout.astro"; -import { api, plural, type Activity, type DictionaryItem } from "../../lib/api"; +import { api, plural, type Activity, type DictionaryItem, type PaginatedActivity } from "../../lib/api"; const { slug } = Astro.params; let fish: DictionaryItem | undefined, items: Activity[] = [], unavailable = false; try { const fishes = await api("/api/v1/fishes?limit=500"); fish = fishes.find(item => item.slug === slug); - if (fish) items = await api(`/api/v1/activity?hours=72&fish=${encodeURIComponent(fish.slug)}&limit=100`); + if (fish) { const paginated = await api(`/api/v1/activity?hours=72&fish=${encodeURIComponent(fish.slug)}&limit=100`); items = paginated.items; } } catch { unavailable = true; } if (unavailable) { Astro.response.status = 503; diff --git a/apps/web/src/pages/spots/[id].astro b/apps/web/src/pages/spots/[id].astro index 0ab50a6..2847f20 100644 --- a/apps/web/src/pages/spots/[id].astro +++ b/apps/web/src/pages/spots/[id].astro @@ -4,7 +4,7 @@ import SourceBadge from "../../components/SourceBadge.astro"; import CoordinateRadar from "../../components/CoordinateRadar.astro"; import ActivityTimeline from "../../components/ActivityTimeline.astro"; import CatchList from "../../components/CatchList.astro"; -import { activityLevel, api, ApiError, plural, type Activity, type Catch, type Spot } from "../../lib/api"; +import { activityLevel, api, ApiError, plural, type Activity, type Catch, type PaginatedActivity, type Spot } from "../../lib/api"; const { id } = Astro.params; let spot: Spot | null = null, catches: Catch[] = [], activity: Activity | null = null, unavailable = false; let timeline: { start: string; end: string; count: number }[] = []; @@ -14,8 +14,8 @@ try { ? await api(`/api/v1/spots/resolve?waterbody=${encodeURIComponent(readable[1])}&x=${readable[2]}&y=${readable[3]}`) : await api(`/api/v1/spots/${id}`); if (!readable) return Astro.redirect(`/spots/${spotResult.waterbody_slug}-${spotResult.x}x${spotResult.y}`, 301); - const [catchResult, activityRows] = await Promise.all([api(`/api/v1/spots/${spotResult.id}/catches`), api(`/api/v1/activity?hours=24&waterbody=${encodeURIComponent(spotResult.waterbody_slug)}&limit=100`)]); - spot = spotResult; catches = catchResult; activity = activityRows.find(item => item.spot_id === spotResult.id) ?? null; + const [catchResult, activityPaginated] = await Promise.all([api(`/api/v1/spots/${spotResult.id}/catches`), api(`/api/v1/activity?hours=24&waterbody=${encodeURIComponent(spotResult.waterbody_slug)}&limit=100`)]); + spot = spotResult; catches = catchResult; activity = activityPaginated.items.find(item => item.spot_id === spotResult.id) ?? null; timeline = await api(`/api/v1/spots/${spotResult.id}/timeline`); } catch (error) { unavailable = true; diff --git a/apps/web/src/pages/waterbodies/[slug].astro b/apps/web/src/pages/waterbodies/[slug].astro index c69d1db..5451582 100644 --- a/apps/web/src/pages/waterbodies/[slug].astro +++ b/apps/web/src/pages/waterbodies/[slug].astro @@ -1,13 +1,13 @@ --- import ActivityCard from "../../components/ActivityCard.astro"; import Layout from "../../layouts/Layout.astro"; -import { api, plural, type Activity, type DictionaryItem } from "../../lib/api"; +import { api, plural, type Activity, type DictionaryItem, type PaginatedActivity } from "../../lib/api"; const { slug } = Astro.params; let water: DictionaryItem | undefined, items: Activity[] = [], unavailable = false; try { const waters = await api("/api/v1/waterbodies?limit=500"); water = waters.find(item => item.slug === slug); - if (water) items = await api(`/api/v1/activity?hours=72&waterbody=${encodeURIComponent(water.slug)}&limit=100`); + if (water) { const paginated = await api(`/api/v1/activity?hours=72&waterbody=${encodeURIComponent(water.slug)}&limit=100`); items = paginated.items; } } catch { unavailable = true; } if (unavailable) { Astro.response.status = 503; diff --git a/apps/web/src/pages/waterbodies/[slug]/[fish].astro b/apps/web/src/pages/waterbodies/[slug]/[fish].astro index 644e09c..6c17369 100644 --- a/apps/web/src/pages/waterbodies/[slug]/[fish].astro +++ b/apps/web/src/pages/waterbodies/[slug]/[fish].astro @@ -1,13 +1,13 @@ --- import ActivityCard from "../../../components/ActivityCard.astro"; import Layout from "../../../layouts/Layout.astro"; -import { api, plural, type Activity, type DictionaryItem } from "../../../lib/api"; +import { api, plural, type Activity, type DictionaryItem, type PaginatedActivity } from "../../../lib/api"; const { slug, fish: fishSlug } = Astro.params; let water: DictionaryItem | undefined, fish: DictionaryItem | undefined, items: Activity[] = [], unavailable = false; try { const [waters, fishes] = await Promise.all([api("/api/v1/waterbodies?limit=500"), api("/api/v1/fishes?limit=500")]); water = waters.find(item => item.slug === slug); fish = fishes.find(item => item.slug === fishSlug); - if (water && fish) items = await api(`/api/v1/activity?hours=72&waterbody=${encodeURIComponent(water.slug)}&fish=${encodeURIComponent(fish.slug)}&limit=100`); + if (water && fish) { const paginated = await api(`/api/v1/activity?hours=72&waterbody=${encodeURIComponent(water.slug)}&fish=${encodeURIComponent(fish.slug)}&limit=100`); items = paginated.items; } } catch { unavailable = true; } if (unavailable) { Astro.response.status = 503;