perf: cache public activity aggregates
This commit is contained in:
@@ -16,6 +16,7 @@ OFFICIAL_RECORDS_CATEGORY=records
|
|||||||
OFFICIAL_IMPORT_REQUIRED=false
|
OFFICIAL_IMPORT_REQUIRED=false
|
||||||
IMPORT_INTERVAL_SECONDS=3600
|
IMPORT_INTERVAL_SECONDS=3600
|
||||||
COMMUNITY_IMPORT_INTERVAL_SECONDS=1800
|
COMMUNITY_IMPORT_INTERVAL_SECONDS=1800
|
||||||
|
PUBLIC_CACHE_SECONDS=20
|
||||||
RF4MAP_POINT_URL=https://rf4map.ru/points/275
|
RF4MAP_POINT_URL=https://rf4map.ru/points/275
|
||||||
RF4POSTS_SPOT_URL=https://rf4-posts.com/ru/spots/d0c6d9c6-4ebf-49a7-98a8-9a562553a8ee
|
RF4POSTS_SPOT_URL=https://rf4-posts.com/ru/spots/d0c6d9c6-4ebf-49a7-98a8-9a562553a8ee
|
||||||
RATE_LIMIT_SECRET=change-rate-limit-secret
|
RATE_LIMIT_SECRET=change-rate-limit-secret
|
||||||
|
|||||||
@@ -10,6 +10,7 @@ POSTGRES_PASSWORD=replace-with-long-random-value
|
|||||||
DATABASE_URL=postgresql+psycopg://rf4:replace-with-url-encoded-password@db:5432/rf4_spotter
|
DATABASE_URL=postgresql+psycopg://rf4:replace-with-url-encoded-password@db:5432/rf4_spotter
|
||||||
APP_VERSION=0.1.0
|
APP_VERSION=0.1.0
|
||||||
APP_REVISION=replace-with-git-commit-sha
|
APP_REVISION=replace-with-git-commit-sha
|
||||||
|
PUBLIC_CACHE_SECONDS=20
|
||||||
|
|
||||||
ADMIN_TOKEN=replace-with-at-least-32-random-characters
|
ADMIN_TOKEN=replace-with-at-least-32-random-characters
|
||||||
RATE_LIMIT_SECRET=replace-with-at-least-32-random-characters
|
RATE_LIMIT_SECRET=replace-with-at-least-32-random-characters
|
||||||
|
|||||||
@@ -88,7 +88,7 @@ docker compose up --build
|
|||||||
|
|
||||||
Контейнер API сам выполняет `alembic upgrade head`, затем идемпотентный seed. PostgreSQL хранит данные в именованном volume `postgres_data`, а MinIO — в `minio_data`. Compose ожидает readiness PostgreSQL и MinIO перед API, а API-контейнер проверяет `/ready`. Версия и commit SHA задаются через `APP_VERSION`/`APP_REVISION`; те же значения доступны администратору в `/api/v1/admin/diagnostics`. Официальный импорт по умолчанию необязателен; при включённом scheduler установите `OFFICIAL_IMPORT_REQUIRED=true`, тогда отсутствующий, неуспешный или просроченный запуск сделает readiness отрицательным.
|
Контейнер API сам выполняет `alembic upgrade head`, затем идемпотентный seed. PostgreSQL хранит данные в именованном volume `postgres_data`, а MinIO — в `minio_data`. Compose ожидает readiness PostgreSQL и MinIO перед API, а API-контейнер проверяет `/ready`. Версия и commit SHA задаются через `APP_VERSION`/`APP_REVISION`; те же значения доступны администратору в `/api/v1/admin/diagnostics`. Официальный импорт по умолчанию необязателен; при включённом scheduler установите `OFFICIAL_IMPORT_REQUIRED=true`, тогда отсутствующий, неуспешный или просроченный запуск сделает readiness отрицательным.
|
||||||
|
|
||||||
API и scheduler пишут по одной JSON-записи на событие. HTTP-лог содержит только сгенерированный `request_id`, метод, путь без query string, статус и длительность; IP, заголовок авторизации и пользовательский payload не журналируются. `X-Request-ID` возвращается клиенту. Стандартный access-log Uvicorn отключён. Уровень управляется `LOG_LEVEL`. Защищённый `/api/v1/admin/diagnostics` скачивает JSON только с идентификатором сборки и агрегированными счётчиками, без имён игроков, исходных URL, payload и ошибок парсеров.
|
API и scheduler пишут по одной JSON-записи на событие. HTTP-лог содержит только сгенерированный `request_id`, метод, путь без query string, статус и длительность; IP, заголовок авторизации и пользовательский payload не журналируются. `X-Request-ID` возвращается клиенту. Стандартный access-log Uvicorn отключён. Уровень управляется `LOG_LEVEL`. Публичный агрегат активности кэшируется в памяти процесса на 20 секунд и очищается после публикации, модерации или удаления; `X-Cache` показывает `HIT`/`MISS`. Защищённый `/api/v1/admin/diagnostics` скачивает JSON только с идентификатором сборки и агрегированными счётчиками, без имён игроков, исходных URL, payload и ошибок парсеров.
|
||||||
|
|
||||||
Остановка:
|
Остановка:
|
||||||
|
|
||||||
|
|||||||
@@ -27,6 +27,7 @@ class Settings(BaseSettings):
|
|||||||
retention_published_payload_days: int = Field(default=365, ge=90)
|
retention_published_payload_days: int = Field(default=365, ge=90)
|
||||||
import_interval_seconds: int = Field(default=3600, ge=3600)
|
import_interval_seconds: int = Field(default=3600, ge=3600)
|
||||||
community_import_interval_seconds: int = Field(default=1800, ge=1800)
|
community_import_interval_seconds: int = Field(default=1800, ge=1800)
|
||||||
|
public_cache_seconds: int = Field(default=20, ge=1, le=300)
|
||||||
rf4map_point_url: str = "https://rf4map.ru/points/275"
|
rf4map_point_url: str = "https://rf4map.ru/points/275"
|
||||||
rf4posts_spot_url: str = "https://rf4-posts.com/ru/spots/d0c6d9c6-4ebf-49a7-98a8-9a562553a8ee"
|
rf4posts_spot_url: str = "https://rf4-posts.com/ru/spots/d0c6d9c6-4ebf-49a7-98a8-9a562553a8ee"
|
||||||
rate_limit_secret: str = "change-rate-limit-secret"
|
rate_limit_secret: str = "change-rate-limit-secret"
|
||||||
|
|||||||
+13
-2
@@ -26,6 +26,7 @@ from .importer import ImportAlreadyRunning, ImportSourceError, import_records, n
|
|||||||
from .logging_config import configure_logging
|
from .logging_config import configure_logging
|
||||||
from .models import Bait, BaitKind, CatchReport, CommunityImportRun, 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 .readiness import readiness_report
|
||||||
|
from .public_cache import public_cache
|
||||||
from .schemas import ActivityOut, AdminCatchReportOut, BaitOut, CatchOut, CatchReportAccepted, CatchReportCreate, CatchReportCreated, ExternalObservationDecision, ExternalObservationMapping, ExternalObservationOut, ExternalObservationPublished, FishOut, ImportRunOut, ModerationUpdate, OfficialRecordOut, PublicObservationOut, SourceStatusOut, 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
|
from .storage import ScreenshotError, client as storage_client, delete_screenshot, signed_screenshot_url, upload_screenshot
|
||||||
|
|
||||||
@@ -110,7 +111,7 @@ def baits(db: Db, limit: int = Query(200, ge=1, le=500), offset: int = Query(0,
|
|||||||
|
|
||||||
@app.get("/api/v1/activity", response_model=list[ActivityOut])
|
@app.get("/api/v1/activity", response_model=list[ActivityOut])
|
||||||
def activity(
|
def activity(
|
||||||
db: Db, hours: int = Query(24),
|
db: Db, response: Response, hours: int = Query(24),
|
||||||
waterbody: str | None = None, fish: str | None = None,
|
waterbody: str | None = None, fish: str | None = None,
|
||||||
method: str | None = None,
|
method: str | None = None,
|
||||||
sort: Literal["activity", "confidence", "freshness"] = "activity",
|
sort: Literal["activity", "confidence", "freshness"] = "activity",
|
||||||
@@ -118,6 +119,12 @@ def activity(
|
|||||||
) -> list[ActivityOut]:
|
) -> list[ActivityOut]:
|
||||||
if hours not in {6, 12, 24, 72}:
|
if hours not in {6, 12, 24, 72}:
|
||||||
raise HTTPException(status_code=422, detail="hours must be one of: 6, 12, 24, 72")
|
raise HTTPException(status_code=422, detail="hours must be one of: 6, 12, 24, 72")
|
||||||
|
response.headers["Cache-Control"] = f"public, max-age={settings.public_cache_seconds}"
|
||||||
|
cache_key = ("activity", hours, waterbody, fish, method, sort, limit, offset)
|
||||||
|
cached = public_cache.get(cache_key, settings.public_cache_seconds)
|
||||||
|
if cached is not None:
|
||||||
|
response.headers["X-Cache"] = "HIT"
|
||||||
|
return cached
|
||||||
rows = activity_rows(db, hours=hours, waterbody=waterbody, fish=fish, method=method)
|
rows = activity_rows(db, hours=hours, waterbody=waterbody, fish=fish, method=method)
|
||||||
keys = {
|
keys = {
|
||||||
"activity": lambda r: (r.activity_score, r.confidence_score, r.last_confirmed_at, str(r.spot_id)),
|
"activity": lambda r: (r.activity_score, r.confidence_score, r.last_confirmed_at, str(r.spot_id)),
|
||||||
@@ -125,7 +132,8 @@ def activity(
|
|||||||
"freshness": lambda r: (r.last_confirmed_at, r.activity_score, r.confidence_score, str(r.spot_id)),
|
"freshness": lambda r: (r.last_confirmed_at, r.activity_score, r.confidence_score, str(r.spot_id)),
|
||||||
}
|
}
|
||||||
rows.sort(key=keys[sort], reverse=True)
|
rows.sort(key=keys[sort], reverse=True)
|
||||||
return rows[offset:offset + limit]
|
response.headers["X-Cache"] = "MISS"
|
||||||
|
return public_cache.set(cache_key, rows[offset:offset + limit])
|
||||||
|
|
||||||
|
|
||||||
def _spot_or_404(db: Session, spot_id: UUID) -> Spot:
|
def _spot_or_404(db: Session, spot_id: UUID) -> Spot:
|
||||||
@@ -372,6 +380,7 @@ def admin_publish_external_observation(
|
|||||||
report = publish_observation(db, observation)
|
report = publish_observation(db, observation)
|
||||||
except ExternalReviewError as exc:
|
except ExternalReviewError as exc:
|
||||||
raise HTTPException(status_code=409, detail=str(exc)) from exc
|
raise HTTPException(status_code=409, detail=str(exc)) from exc
|
||||||
|
public_cache.invalidate()
|
||||||
return ExternalObservationPublished(observation_id=observation.id, catch_report_id=report.id, status=observation.status)
|
return ExternalObservationPublished(observation_id=observation.id, catch_report_id=report.id, status=observation.status)
|
||||||
|
|
||||||
|
|
||||||
@@ -454,6 +463,7 @@ def moderate_report(report_id: UUID, payload: ModerationUpdate, db: Db, moderato
|
|||||||
report.moderation_status = ModerationStatus(payload.status)
|
report.moderation_status = ModerationStatus(payload.status)
|
||||||
db.add(ModerationEvent(catch_report=report, created_at=datetime.now(timezone.utc), previous_status=previous, new_status=report.moderation_status, moderator=moderator, reason=payload.reason))
|
db.add(ModerationEvent(catch_report=report, created_at=datetime.now(timezone.utc), previous_status=previous, new_status=report.moderation_status, moderator=moderator, reason=payload.reason))
|
||||||
db.commit()
|
db.commit()
|
||||||
|
public_cache.invalidate()
|
||||||
return CatchReportCreated(id=report.id, moderation_status=report.moderation_status.value)
|
return CatchReportCreated(id=report.id, moderation_status=report.moderation_status.value)
|
||||||
|
|
||||||
|
|
||||||
@@ -476,6 +486,7 @@ def delete_report(report_id: UUID, db: Db, moderator: Annotated[str, Depends(_ad
|
|||||||
report.raw_payload = None
|
report.raw_payload = None
|
||||||
db.add(ModerationEvent(catch_report=report, created_at=report.deleted_at, previous_status=previous, new_status=ModerationStatus.rejected, moderator=moderator, reason="user report deleted and anonymized"))
|
db.add(ModerationEvent(catch_report=report, created_at=report.deleted_at, previous_status=previous, new_status=ModerationStatus.rejected, moderator=moderator, reason="user report deleted and anonymized"))
|
||||||
db.commit()
|
db.commit()
|
||||||
|
public_cache.invalidate()
|
||||||
return Response(status_code=204)
|
return Response(status_code=204)
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@@ -0,0 +1,32 @@
|
|||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
from copy import deepcopy
|
||||||
|
from threading import Lock
|
||||||
|
from time import monotonic
|
||||||
|
from typing import Any
|
||||||
|
|
||||||
|
|
||||||
|
class PublicResponseCache:
|
||||||
|
def __init__(self) -> None:
|
||||||
|
self._items: dict[tuple[Any, ...], tuple[float, Any]] = {}
|
||||||
|
self._lock = Lock()
|
||||||
|
|
||||||
|
def get(self, key: tuple[Any, ...], ttl_seconds: int) -> Any | None:
|
||||||
|
with self._lock:
|
||||||
|
item = self._items.get(key)
|
||||||
|
if item is None or monotonic() - item[0] >= ttl_seconds:
|
||||||
|
self._items.pop(key, None)
|
||||||
|
return None
|
||||||
|
return deepcopy(item[1])
|
||||||
|
|
||||||
|
def set(self, key: tuple[Any, ...], value: Any) -> Any:
|
||||||
|
with self._lock:
|
||||||
|
self._items[key] = (monotonic(), deepcopy(value))
|
||||||
|
return value
|
||||||
|
|
||||||
|
def invalidate(self) -> None:
|
||||||
|
with self._lock:
|
||||||
|
self._items.clear()
|
||||||
|
|
||||||
|
|
||||||
|
public_cache = PublicResponseCache()
|
||||||
@@ -0,0 +1,11 @@
|
|||||||
|
from app.public_cache import PublicResponseCache
|
||||||
|
|
||||||
|
|
||||||
|
def test_cache_copies_values_and_invalidates() -> None:
|
||||||
|
cache = PublicResponseCache()
|
||||||
|
original = [{"score": 10}]
|
||||||
|
cache.set(("activity",), original)
|
||||||
|
original[0]["score"] = 99
|
||||||
|
assert cache.get(("activity",), 20) == [{"score": 10}]
|
||||||
|
cache.invalidate()
|
||||||
|
assert cache.get(("activity",), 20) is None
|
||||||
@@ -121,6 +121,7 @@ services:
|
|||||||
OFFICIAL_IMPORT_REQUIRED: ${OFFICIAL_IMPORT_REQUIRED:-false}
|
OFFICIAL_IMPORT_REQUIRED: ${OFFICIAL_IMPORT_REQUIRED:-false}
|
||||||
SEED_DEMO_DATA: "false"
|
SEED_DEMO_DATA: "false"
|
||||||
IMPORT_INTERVAL_SECONDS: ${IMPORT_INTERVAL_SECONDS:-3600}
|
IMPORT_INTERVAL_SECONDS: ${IMPORT_INTERVAL_SECONDS:-3600}
|
||||||
|
PUBLIC_CACHE_SECONDS: ${PUBLIC_CACHE_SECONDS:-20}
|
||||||
RATE_LIMIT_SECRET: ${RATE_LIMIT_SECRET:?Set RATE_LIMIT_SECRET}
|
RATE_LIMIT_SECRET: ${RATE_LIMIT_SECRET:?Set RATE_LIMIT_SECRET}
|
||||||
RETENTION_SUBMISSION_DAYS: ${RETENTION_SUBMISSION_DAYS:-1}
|
RETENTION_SUBMISSION_DAYS: ${RETENTION_SUBMISSION_DAYS:-1}
|
||||||
RETENTION_UNREVIEWED_DAYS: ${RETENTION_UNREVIEWED_DAYS:-30}
|
RETENTION_UNREVIEWED_DAYS: ${RETENTION_UNREVIEWED_DAYS:-30}
|
||||||
|
|||||||
@@ -33,6 +33,8 @@ docker run --rm caddy:2.10.2-alpine caddy hash-password --plaintext 'ОТДЕЛ
|
|||||||
|
|
||||||
Перед сборкой запишите текущий `git rev-parse --short HEAD` в `APP_REVISION` файла `.env.production`, чтобы `/ready` однозначно показывал развёрнутый commit.
|
Перед сборкой запишите текущий `git rev-parse --short HEAD` в `APP_REVISION` файла `.env.production`, чтобы `/ready` однозначно показывал развёрнутый commit.
|
||||||
|
|
||||||
|
`PUBLIC_CACHE_SECONDS` задаёт короткий in-process кэш публичного агрегата активности (по умолчанию 20 секунд). Не увеличивайте его без повторной проверки свежести после публикации.
|
||||||
|
|
||||||
```bash
|
```bash
|
||||||
docker compose --env-file .env.production -f compose.production.yaml config --quiet
|
docker compose --env-file .env.production -f compose.production.yaml config --quiet
|
||||||
docker compose --env-file .env.production -f compose.production.yaml build
|
docker compose --env-file .env.production -f compose.production.yaml build
|
||||||
|
|||||||
+1
-1
@@ -141,7 +141,7 @@
|
|||||||
- [x] Добавить scheduler всех пяти разрешённых community-парсеров с устойчивым cooldown ≥30 минут, PostgreSQL lock, журналом запусков и opt-in локальным профилем; backoff после повторных ошибок остаётся отдельным улучшением (7 сентября 2026).
|
- [x] Добавить scheduler всех пяти разрешённых community-парсеров с устойчивым cooldown ≥30 минут, PostgreSQL lock, журналом запусков и opt-in локальным профилем; backoff после повторных ошибок остаётся отдельным улучшением (7 сентября 2026).
|
||||||
- [x] Добавить экспоненциальный backoff community scheduler после повторных ошибок от 30 минут до 24 часов со сбросом после успеха (7 сентября 2026).
|
- [x] Добавить экспоненциальный backoff community scheduler после повторных ошибок от 30 минут до 24 часов со сбросом после успеха (7 сентября 2026).
|
||||||
- [x] Показывать на `/status` безопасное состояние источников: актуален, устарел, временно ограничен, изменился DOM, ожидает запуска или выключен (7 сентября 2026).
|
- [x] Показывать на `/status` безопасное состояние источников: актуален, устарел, временно ограничен, изменился DOM, ожидает запуска или выключен (7 сентября 2026).
|
||||||
- [ ] Добавить короткое серверное кэширование публичных GET API и проверить корректную инвалидацию после публикации.
|
- [x] Добавить 20-секундный серверный кэш агрегата активности с явной инвалидацией после публикации, модерации и удаления (7 сентября 2026).
|
||||||
- [x] Настроить долгий immutable cache для хешированных assets и разумный cache для изображений/favicon (7 сентября 2026).
|
- [x] Настроить долгий immutable cache для хешированных assets и разумный cache для изображений/favicon (7 сентября 2026).
|
||||||
- [x] Добавить серверное «Показать ещё» для публичной ленты полевых сигналов с сохранением фильтров и пределом 48 записей (7 сентября 2026).
|
- [x] Добавить серверное «Показать ещё» для публичной ленты полевых сигналов с сохранением фильтров и пределом 48 записей (7 сентября 2026).
|
||||||
- [x] Объединять одинаковые полевые сигналы в сюжеты без потери уникальных ссылок provenance и показывать число совпадений (7 сентября 2026).
|
- [x] Объединять одинаковые полевые сигналы в сюжеты без потери уникальных ссылок provenance и показывать число совпадений (7 сентября 2026).
|
||||||
|
|||||||
Reference in New Issue
Block a user