diff --git a/.env.example b/.env.example index f0bd407..0d7a4b9 100644 --- a/.env.example +++ b/.env.example @@ -54,6 +54,17 @@ QUERY_LIST_CACHE_MAX_ENTRIES=4 GLOBAL_TREND_CACHE_MAX_ENTRIES=64 QUERY_LIST_CACHE_MAX_ROWS=100000 QUERY_LIST_CACHE_MAX_BYTES=67108864 +# Kalici query-metrics materialized snapshotlarini ayri worker yeniler. API +# request yolunda advisor.query_metrics() calismaz. Concurrent refresh boyunca +# onceki tamamlanmis snapshot servis edilmeye devam eder. +QUERY_METRICS_SNAPSHOT_POLL_SECONDS=15 +QUERY_METRICS_SNAPSHOT_1H_REFRESH_SECONDS=900 +QUERY_METRICS_SNAPSHOT_24H_REFRESH_SECONDS=3600 +QUERY_METRICS_SNAPSHOT_7D_REFRESH_SECONDS=21600 +QUERY_METRICS_SNAPSHOT_30D_REFRESH_SECONDS=43200 +QUERY_METRICS_SNAPSHOT_STATEMENT_TIMEOUT_SECONDS=1800 +QUERY_METRICS_SNAPSHOT_RETRY_SECONDS=60 +QUERY_METRICS_SNAPSHOT_WORKER_MEMORY_LIMIT=256m # Tam snapshotlarin Python nesne ek yükü ayrıca container seviyesinde sınırlıdır. API_MEMORY_LIMIT=1g # Eszamanli repository migration runner'larinin advisory lock bekleme siniri. @@ -91,6 +102,10 @@ WORKLOAD_PROFILE=normal WORKLOAD_DURATION_SECONDS=0 WORKLOAD_WORKERS=6 WORKLOAD_INTERVAL_SECONDS=0.25 +WORKLOAD_INTERVAL_JITTER_RATIO=0 +WORKLOAD_TRAFFIC_PHASE_SECONDS=0 +WORKLOAD_TRAFFIC_MIN_INTERVAL_MULTIPLIER=1 +WORKLOAD_TRAFFIC_MAX_INTERVAL_MULTIPLIER=1 WORKLOAD_REPORT_INTERVAL_SECONDS=10 WORKLOAD_RANDOM_SEED=20260725 WORKLOAD_STATEMENT_TIMEOUT_MS=15000 diff --git a/.github/workflows/postgres-integration.yml b/.github/workflows/postgres-integration.yml index dc82e39..9404f83 100644 --- a/.github/workflows/postgres-integration.yml +++ b/.github/workflows/postgres-integration.yml @@ -95,15 +95,18 @@ jobs: - name: Verify health and migration idempotency run: | - curl --fail --silent --show-error \ - http://127.0.0.1:8000/api/v1/health \ - | python -c ' - import json, sys - health = json.load(sys.stdin) - assert health["status"] == "healthy", health - assert health["repository"] == "healthy", health - assert health["collector"] == "healthy", health - ' + health_ready=false + for attempt in $(seq 1 18); do + if curl --fail --silent --show-error \ + http://127.0.0.1:8000/api/v1/health \ + | python -c 'import json, sys; health = json.load(sys.stdin); assert health["status"] == "healthy", health; assert health["repository"] == "healthy", health; assert health["collector"] == "healthy", health'; then + health_ready=true + break + fi + printf 'collector health attempt=%s is not ready yet\n' "$attempt" + sleep 10 + done + test "$health_ready" = true docker compose run --rm repository-migrate bash scripts/verify-temporal-reliability.sh docker compose exec -T repository-db \ @@ -118,7 +121,7 @@ jobs: psql -X --set=ON_ERROR_STOP=1 --username postgres --dbname powa --file=- \ < sql/tests/join_outbox_guardrail_integration.sql - - name: Verify a persisted 0013 database upgrades to 0014 + - name: Verify a persisted 0013 database upgrades through 0016 shell: bash run: | set -Eeuo pipefail @@ -170,7 +173,7 @@ jobs: release record; BEGIN SELECT * INTO STRICT release FROM advisor.release_info(); - IF release.current_migration <> '0014' OR release.applied_count <> 14 THEN + IF release.current_migration <> '0016' OR release.applied_count <> 16 THEN RAISE EXCEPTION 'unexpected upgraded release state: %', row_to_json(release); END IF; IF NOT has_function_privilege( @@ -207,6 +210,27 @@ jobs: - name: Prepare the deterministic runtime candidate run: bash scripts/verify.sh + - name: Refresh the disposable 30d runtime replay snapshot + run: | + docker compose exec -T query-metrics-snapshot-worker python - <<'PY' + import time + + from app.config import get_settings + from app.snapshot_worker import open_connection, refresh_snapshot + + + settings = get_settings() + with open_connection(settings) as connection: + for _ in range(30): + if refresh_snapshot(connection, "30d"): + break + time.sleep(2) + else: + raise SystemExit( + "30d query metrics snapshot advisory lock could not be acquired" + ) + PY + - name: Boot disposable clone services run: | docker compose --profile real-validation up -d --wait \ diff --git a/CHANGELOG.md b/CHANGELOG.md index 20e6a04..b7949ee 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -4,6 +4,15 @@ Bu proje [Semantic Versioning](https://semver.org/) kullanır. Repository migration'ları ileri yönlüdür; uygulama sürümü ile şema uyumluluğu yükseltme runbook'unda belirtilir. +## 1.1.1 — 2026-07-30 + +- 1h/24h/7d/30d query-metrics sonuçları kalıcı materialized snapshot olarak + request yolundan önce hesaplanır. +- Ayrı worker snapshot'ları sırayla ve atomik olarak yeniler; hesaplama sürerken + önceki tamamlanmış sonuç servis edilir. +- Repository migration hedefi `0016`; query-metrics snapshot'larına ek olarak + global/server/database overview trendleri de önceden hesaplanır. + ## 1.1.0 — 2026-07-26 - Ana kaynakta yalnız persisted, salt-okunur sorgular için gerçek `EXPLAIN ANALYZE`. diff --git a/README.md b/README.md index 9a3f2be..88286db 100644 --- a/README.md +++ b/README.md @@ -1,6 +1,6 @@ # PostgreSQL Sorgu Performansı ve Öneri Motoru -Güncel sürüm: `1.1.0`. Değişiklikler [CHANGELOG.md](CHANGELOG.md), güvenli +Güncel sürüm: `1.1.1`. Değişiklikler [CHANGELOG.md](CHANGELOG.md), güvenli yükseltme ve geri dönüş adımları [upgrade/rollback runbook'unda](docs/UPGRADE_ROLLBACK.md). PDF v1.1'de tarif edilen ilk iterasyonun çalışan referans uygulamasıdır. Tek bir Docker/OrbStack hostu üzerinde **iki ayrı PostgreSQL sunucu süreci** çalışır: demo kaynak instance `5432`, PoWA repository instance `5433`. PoWA Collector istatistikleri kaynaktan repository'ye taşır; FastAPI yalnız repository'yi okur ve React arayüzü sonuçları gösterir. Aynı repository/collector, `scripts/register-source.sh` ile birden fazla gerçek PostgreSQL kaynağı izleyebilir. diff --git a/VERSION b/VERSION index 9084fa2..524cb55 100644 --- a/VERSION +++ b/VERSION @@ -1 +1 @@ -1.1.0 +1.1.1 diff --git a/backend/app/config.py b/backend/app/config.py index a8292b6..e2b0496 100644 --- a/backend/app/config.py +++ b/backend/app/config.py @@ -24,6 +24,20 @@ "30d": "1 day", } +QUERY_METRICS_SNAPSHOT_VIEWS: dict[str, str] = { + "1h": "query_metrics_snapshot_1h", + "24h": "query_metrics_snapshot_24h", + "7d": "query_metrics_snapshot_7d", + "30d": "query_metrics_snapshot_30d", +} + +GLOBAL_TREND_SNAPSHOT_VIEWS: dict[str, str] = { + "1h": "global_trend_snapshot_1h", + "24h": "global_trend_snapshot_24h", + "7d": "global_trend_snapshot_7d", + "30d": "global_trend_snapshot_30d", +} + PrincipalRole = Literal["analyst", "annotator", "admin"] @@ -91,6 +105,26 @@ class Settings(BaseSettings): ge=1024 * 1024, le=1024 * 1024 * 1024, ) + # A dedicated worker refreshes persistent materialized snapshots. API + # requests only read those snapshots and therefore never execute the + # expensive query_metrics function on a cold process-local cache. + query_metrics_snapshot_poll_seconds: float = Field(default=15.0, ge=1, le=300) + query_metrics_snapshot_1h_refresh_seconds: int = Field( + default=15 * 60, ge=60, le=86_400 + ) + query_metrics_snapshot_24h_refresh_seconds: int = Field( + default=60 * 60, ge=60, le=7 * 86_400 + ) + query_metrics_snapshot_7d_refresh_seconds: int = Field( + default=6 * 60 * 60, ge=60, le=30 * 86_400 + ) + query_metrics_snapshot_30d_refresh_seconds: int = Field( + default=12 * 60 * 60, ge=60, le=30 * 86_400 + ) + query_metrics_snapshot_statement_timeout_seconds: int = Field( + default=30 * 60, ge=60, le=6 * 60 * 60 + ) + query_metrics_snapshot_retry_seconds: int = Field(default=60, ge=5, le=3_600) sql_text_visibility: str = "authorized" retention_days: int = 90 log_level: str = "INFO" diff --git a/backend/app/main.py b/backend/app/main.py index 67019ad..f301f86 100644 --- a/backend/app/main.py +++ b/backend/app/main.py @@ -7,12 +7,15 @@ from fastapi import FastAPI from fastapi.middleware.cors import CORSMiddleware from prometheus_client import CONTENT_TYPE_LATEST, generate_latest +from starlette import status +from starlette.requests import Request +from starlette.responses import JSONResponse from starlette.responses import Response from app.api.router import router from app.config import get_settings from app.db import close_pool, open_pool -from app.repositories.powa import repository +from app.repositories.powa import QueryMetricsSnapshotWarming, repository from app.version import APPLICATION_VERSION @@ -58,6 +61,18 @@ async def lifespan(_: FastAPI): app.include_router(router) +@app.exception_handler(QueryMetricsSnapshotWarming) +async def snapshot_warming( + _: Request, + __: QueryMetricsSnapshotWarming, +) -> JSONResponse: + return JSONResponse( + status_code=status.HTTP_503_SERVICE_UNAVAILABLE, + content={"detail": "Dashboard verileri ilk kez hazirlaniyor; kisa sure sonra deneyin."}, + headers={"Retry-After": "30"}, + ) + + @app.get("/", include_in_schema=False) async def root() -> dict[str, str]: return {"name": app.title, "docs": "/docs", "health": "/api/v1/health"} diff --git a/backend/app/repositories/powa.py b/backend/app/repositories/powa.py index 215d2ca..e50e2d3 100644 --- a/backend/app/repositories/powa.py +++ b/backend/app/repositories/powa.py @@ -13,7 +13,15 @@ from typing import Any from uuid import NAMESPACE_URL, uuid5 -from app.config import WINDOW_BUCKETS, WINDOW_INTERVALS, get_settings +from psycopg.errors import ObjectNotInPrerequisiteState + +from app.config import ( + GLOBAL_TREND_SNAPSHOT_VIEWS, + QUERY_METRICS_SNAPSHOT_VIEWS, + WINDOW_BUCKETS, + WINDOW_INTERVALS, + get_settings, +) from app.db import pool @@ -44,6 +52,14 @@ class QueryMetricsRefreshBackoff(RuntimeError): """A failed refresh is inside its bounded repository-protection backoff.""" +class QueryMetricsSnapshotWarming(RuntimeError): + """The persistent snapshot has not completed its first refresh yet.""" + + +class GlobalTrendSnapshotWarming(QueryMetricsSnapshotWarming): + """The persistent overview trend has not completed its first refresh yet.""" + + class GlobalTrendRefreshBackoff(RuntimeError): """A failed global-trend refresh is inside its bounded retry backoff.""" @@ -580,12 +596,11 @@ def __init__( self._query_metrics_cache_generation = 0 self._query_metrics_retry: dict[str, _QueryMetricsRetryState] = {} self._query_metrics_cache_lock = asyncio.Lock() - # Deliberately shared by query_metrics and global query_trend across - # every window. The repository must never run two expensive dashboard - # refresh scans concurrently. + # Global trends still require expensive live telemetry scans and share + # this lock across windows. Query-metrics reloads read precomputed + # materialized snapshots and intentionally do not queue behind it. self._repository_refresh_lock = asyncio.Lock() - # Retain the old private name for compatibility with diagnostics while - # making its repository-wide scope explicit above. + # Retain the old private name for compatibility with diagnostics. self._query_metrics_refresh_lock = self._repository_refresh_lock self._query_metrics_refresh_tasks: dict[ tuple[int, str], asyncio.Task[list[dict[str, Any]]] @@ -604,39 +619,47 @@ def __init__( self._clock = clock async def _load_query_metrics_snapshot(self, window: str) -> list[dict[str, Any]]: - interval = interval_for(window) + interval_for(window) + view_name = QUERY_METRICS_SNAPSHOT_VIEWS[window] async with pool.connection() as connection: - # A transaction-local application name remains visible while the - # named cursor executes FETCH statements. This lets the benchmark - # include stale-while-revalidate work in its measurement boundary. + # This is a bounded read of an already-computed materialized + # snapshot. Live annotations are overlaid so an annotation mutation + # does not have to trigger another full telemetry computation. async with connection.cursor() as control_cursor: await control_cursor.execute( "SET LOCAL application_name = " - "'advisor-query-metrics-cache-refresh'" + "'advisor-query-metrics-snapshot-read'" ) - # LIMIT max+1 detects an oversized row set without truncating it - # silently. The named cursor keeps libpq from buffering the whole - # result, while pg_column_size enforces a second payload-byte cap. async with connection.cursor( - name="advisor_query_metrics_cache_refresh" + name="advisor_query_metrics_snapshot_read" ) as cursor: - await cursor.execute( - """ - /* advisor-query-metrics-cache-refresh */ + try: + await cursor.execute( + f""" + /* advisor-query-metrics-snapshot-read */ SELECT metrics.*, - servers.alias AS server_alias, + COALESCE(annotation.status, 'NEW') + AS _live_review_status, + annotation.note AS _live_note, + annotation.updated_by AS _live_updated_by, + annotation.updated_at AS _live_updated_at, pg_column_size(metrics) - + COALESCE(pg_column_size(servers.alias), 0) AS _cache_row_bytes - FROM advisor.query_metrics(%s::interval) AS metrics - LEFT JOIN "PoWA".powa_servers AS servers - ON servers.id = metrics.server_id + FROM advisor.{view_name} AS metrics + LEFT JOIN advisor.query_annotations AS annotation + ON annotation.server_id = metrics.server_id + AND annotation.database_id = metrics.database_id + AND annotation.query_id = metrics.query_id LIMIT %s - /* advisor-query-metrics-cache-refresh */ + /* advisor-query-metrics-snapshot-read */ """, - (interval, self._query_list_cache_max_rows + 1), - ) + (self._query_list_cache_max_rows + 1,), + ) + except ObjectNotInPrerequisiteState as error: + raise QueryMetricsSnapshotWarming( + f"query metrics snapshot {window} is warming" + ) from error rows: list[dict[str, Any]] = [] payload_bytes = 0 while True: @@ -658,6 +681,16 @@ async def _load_query_metrics_snapshot(self, window: str) -> list[dict[str, Any] raise QueryMetricsSnapshotTooLarge( "query metrics snapshot configured byte limitini asti" ) + row["review_status"] = row.pop( + "_live_review_status", row.get("review_status") or "NEW" + ) + row["note"] = row.pop("_live_note", row.get("note")) + row["updated_by"] = row.pop( + "_live_updated_by", row.get("updated_by") + ) + row["updated_at"] = row.pop( + "_live_updated_at", row.get("updated_at") + ) rows.append(row) if len(rows) > self._query_list_cache_max_rows: raise QueryMetricsSnapshotTooLarge( @@ -710,8 +743,7 @@ async def _refresh_query_metrics_snapshot( ) -> list[dict[str, Any]]: task_key = (generation, window) try: - async with self._repository_refresh_lock: - loaded_rows = await self._load_query_metrics_snapshot(window) + loaded_rows = await self._load_query_metrics_snapshot(window) if len(loaded_rows) > self._query_list_cache_max_rows: # Defense in depth for tests/custom repository subclasses that # override the bounded server-cursor loader. @@ -813,49 +845,40 @@ async def _load_global_trend_snapshot( server_id: int | None = None, database_id: int | None = None, ) -> list[dict[str, Any]]: - interval = interval_for(window) - bucket = WINDOW_BUCKETS[window] - if server_id is None: - function_sql = "advisor.query_trend(now() - %s::interval, %s::interval)" - params: list[Any] = [ - interval, - bucket, - self._query_list_cache_max_rows + 1, - ] - else: - function_sql = ( - "advisor.query_trend(" - "now() - %s::interval, %s::interval, %s, %s)" - ) - params = [ - interval, - bucket, - server_id, - database_id, - self._query_list_cache_max_rows + 1, - ] + interval_for(window) + view_name = GLOBAL_TREND_SNAPSHOT_VIEWS[window] async with pool.connection() as connection: - # Keep background SWR work visible to the benchmark so its - # repository cost cannot leak beyond the measurement boundary. async with connection.cursor() as control_cursor: await control_cursor.execute( "SET LOCAL application_name = " - "'advisor-global-trend-cache-refresh'", + "'advisor-global-trend-snapshot-read'", [], ) async with connection.cursor() as cursor: - await cursor.execute( - f""" - /* advisor-global-trend-cache-refresh */ - SELECT bucket_at AS timestamp, + try: + await cursor.execute( + f""" + /* advisor-global-trend-snapshot-read */ + SELECT timestamp, total_exec_time_ms, calls - FROM {function_sql} + FROM advisor.{view_name} + WHERE server_id IS NOT DISTINCT FROM %s + AND database_id IS NOT DISTINCT FROM %s + ORDER BY timestamp LIMIT %s - /* advisor-global-trend-cache-refresh */ + /* advisor-global-trend-snapshot-read */ """, - params, - ) + ( + server_id, + database_id, + self._query_list_cache_max_rows + 1, + ), + ) + except ObjectNotInPrerequisiteState as error: + raise GlobalTrendSnapshotWarming( + f"global trend snapshot {window} is warming" + ) from error rows = [dict(row) for row in await cursor.fetchall()] if len(rows) > self._query_list_cache_max_rows: raise GlobalTrendSnapshotTooLarge( @@ -923,15 +946,14 @@ async def _refresh_global_trend_snapshot( ) task_key = (generation, cache_key) try: - async with self._repository_refresh_lock: - if server_id is None: - loaded_rows = await self._load_global_trend_snapshot(window) - else: - loaded_rows = await self._load_global_trend_snapshot( - window, - server_id=server_id, - database_id=database_id, - ) + if server_id is None: + loaded_rows = await self._load_global_trend_snapshot(window) + else: + loaded_rows = await self._load_global_trend_snapshot( + window, + server_id=server_id, + database_id=database_id, + ) if len(loaded_rows) > self._query_list_cache_max_rows: # Defense in depth for tests/custom subclasses that override # the bounded loader. diff --git a/backend/app/snapshot_worker.py b/backend/app/snapshot_worker.py new file mode 100644 index 0000000..62c697c --- /dev/null +++ b/backend/app/snapshot_worker.py @@ -0,0 +1,369 @@ +from __future__ import annotations + +import logging +import signal +from datetime import datetime, timezone +from threading import Event +from time import monotonic +from typing import Final + +import psycopg +from psycopg import Connection, sql +from psycopg.rows import dict_row + +from app.config import ( + GLOBAL_TREND_SNAPSHOT_VIEWS, + QUERY_METRICS_SNAPSHOT_VIEWS, + Settings, + get_settings, +) + + +logger = logging.getLogger("advisor.snapshot_worker") + +WINDOW_ORDER: Final[tuple[str, ...]] = ("1h", "24h", "7d", "30d") +QUERY_METRICS_STATE_TABLE: Final = "query_metrics_snapshot_state" +GLOBAL_TREND_STATE_TABLE: Final = "global_trend_snapshot_state" +QUERY_METRICS_LOCK_NAME: Final = ( + "postgresql-advisor:query-metrics-snapshot-refresh" +) +GLOBAL_TREND_LOCK_NAME: Final = "postgresql-advisor:global-trend-snapshot-refresh" + + +def refresh_intervals(settings: Settings) -> dict[str, int]: + return { + "1h": settings.query_metrics_snapshot_1h_refresh_seconds, + "24h": settings.query_metrics_snapshot_24h_refresh_seconds, + "7d": settings.query_metrics_snapshot_7d_refresh_seconds, + "30d": settings.query_metrics_snapshot_30d_refresh_seconds, + } + + +def open_connection(settings: Settings) -> Connection[dict[str, object]]: + connection = psycopg.connect( + settings.database_conninfo, + autocommit=True, + row_factory=dict_row, + application_name="advisor-dashboard-snapshot-worker", + ) + connection.execute( + "SELECT set_config('statement_timeout', %s, false)", + (f"{settings.query_metrics_snapshot_statement_timeout_seconds}s",), + ) + return connection + + +def _set_refreshing( + connection: Connection[dict[str, object]], + window: str, + *, + state_table: str, +) -> None: + connection.execute( + f""" + UPDATE advisor.{state_table} + SET status = 'refreshing', + refresh_started_at = clock_timestamp(), + last_error = NULL, + updated_at = clock_timestamp() + WHERE window_key = %s + """, + (window,), + ) + + +def _set_ready( + connection: Connection[dict[str, object]], + window: str, + *, + state_table: str, + duration_ms: int, + row_count: int, +) -> None: + connection.execute( + f""" + UPDATE advisor.{state_table} + SET status = 'ready', + refreshed_at = clock_timestamp(), + refresh_duration_ms = %s, + row_count = %s, + last_error = NULL, + updated_at = clock_timestamp() + WHERE window_key = %s + """, + (duration_ms, row_count, window), + ) + + +def _set_failed( + connection: Connection[dict[str, object]], + window: str, + error: BaseException, + *, + state_table: str, +) -> None: + sqlstate = getattr(error, "sqlstate", None) + diagnostic = type(error).__name__ + if sqlstate: + diagnostic = f"{diagnostic} (SQLSTATE {sqlstate})" + connection.execute( + f""" + UPDATE advisor.{state_table} + SET status = 'failed', + last_error = %s, + updated_at = clock_timestamp() + WHERE window_key = %s + """, + (diagnostic[:500], window), + ) + + +def _refresh_snapshot( + connection: Connection[dict[str, object]], + window: str, + *, + views: dict[str, str], + state_table: str, + lock_name: str, + family: str, +) -> bool: + view_name = views[window] + lock_row = connection.execute( + "SELECT pg_try_advisory_lock(hashtextextended(%s, 0)) AS acquired", + (lock_name,), + ).fetchone() + if not lock_row or not lock_row["acquired"]: + return False + + started_at = monotonic() + try: + _set_refreshing(connection, window, state_table=state_table) + populated_row = connection.execute( + """ + SELECT ispopulated + FROM pg_matviews + WHERE schemaname = 'advisor' + AND matviewname = %s + """, + (view_name,), + ).fetchone() + if populated_row is None: + raise RuntimeError(f"{family} materialized snapshot is missing") + + concurrently = ( + sql.SQL("CONCURRENTLY ") if populated_row["ispopulated"] else sql.SQL("") + ) + connection.execute( + sql.SQL("REFRESH MATERIALIZED VIEW {}{}.{}").format( + concurrently, + sql.Identifier("advisor"), + sql.Identifier(view_name), + ) + ) + count_row = connection.execute( + sql.SQL("SELECT count(*) AS row_count FROM {}.{}").format( + sql.Identifier("advisor"), + sql.Identifier(view_name), + ) + ).fetchone() + duration_ms = max(0, round((monotonic() - started_at) * 1000)) + _set_ready( + connection, + window, + state_table=state_table, + duration_ms=duration_ms, + row_count=int(count_row["row_count"]) if count_row else 0, + ) + logger.info( + "snapshot ready family=%s window=%s rows=%s duration_ms=%s", + family, + window, + count_row["row_count"] if count_row else 0, + duration_ms, + ) + return True + except BaseException as error: + try: + _set_failed( + connection, + window, + error, + state_table=state_table, + ) + except Exception: + logger.exception( + "could not persist snapshot failure family=%s window=%s", + family, + window, + ) + raise + finally: + try: + connection.execute( + "SELECT pg_advisory_unlock(hashtextextended(%s, 0))", + (lock_name,), + ) + except Exception: + logger.exception("could not release snapshot advisory lock") + + +def refresh_snapshot( + connection: Connection[dict[str, object]], + window: str, +) -> bool: + return _refresh_snapshot( + connection, + window, + views=QUERY_METRICS_SNAPSHOT_VIEWS, + state_table=QUERY_METRICS_STATE_TABLE, + lock_name=QUERY_METRICS_LOCK_NAME, + family="query_metrics", + ) + + +def refresh_global_trend_snapshot( + connection: Connection[dict[str, object]], + window: str, +) -> bool: + return _refresh_snapshot( + connection, + window, + views=GLOBAL_TREND_SNAPSHOT_VIEWS, + state_table=GLOBAL_TREND_STATE_TABLE, + lock_name=GLOBAL_TREND_LOCK_NAME, + family="global_trend", + ) + + +def _due_windows( + connection: Connection[dict[str, object]], + settings: Settings, + *, + state_table: str, + now: datetime | None = None, +) -> list[str]: + rows = connection.execute( + f""" + SELECT window_key, refreshed_at + FROM advisor.{state_table} + ORDER BY CASE window_key + WHEN '1h' THEN 1 + WHEN '24h' THEN 2 + WHEN '7d' THEN 3 + WHEN '30d' THEN 4 + END + """ + ).fetchall() + refreshed = {str(row["window_key"]): row["refreshed_at"] for row in rows} + observed_at = now or datetime.now(timezone.utc) + intervals = refresh_intervals(settings) + due: list[str] = [] + for window in WINDOW_ORDER: + refreshed_at = refreshed.get(window) + if not isinstance(refreshed_at, datetime): + due.append(window) + continue + age_seconds = max(0.0, (observed_at - refreshed_at).total_seconds()) + if age_seconds >= intervals[window]: + due.append(window) + return due + + +def due_windows( + connection: Connection[dict[str, object]], + settings: Settings, + *, + now: datetime | None = None, +) -> list[str]: + return _due_windows( + connection, + settings, + state_table=QUERY_METRICS_STATE_TABLE, + now=now, + ) + + +def global_trend_due_windows( + connection: Connection[dict[str, object]], + settings: Settings, + *, + now: datetime | None = None, +) -> list[str]: + return _due_windows( + connection, + settings, + state_table=GLOBAL_TREND_STATE_TABLE, + now=now, + ) + + +def run_worker(settings: Settings, stop: Event) -> None: + retry_not_before: dict[tuple[str, str], float] = {} + while not stop.is_set(): + try: + with open_connection(settings) as connection: + while not stop.is_set(): + jobs = [ + ("query_metrics", window, refresh_snapshot) + for window in due_windows(connection, settings) + ] + jobs.extend( + ( + "global_trend", + window, + refresh_global_trend_snapshot, + ) + for window in global_trend_due_windows(connection, settings) + ) + attempted = False + for family, window, refresh in jobs: + if stop.is_set(): + break + retry_key = (family, window) + if monotonic() < retry_not_before.get(retry_key, 0): + continue + attempted = True + try: + refreshed = refresh(connection, window) + if refreshed: + retry_not_before.pop(retry_key, None) + else: + retry_not_before[retry_key] = ( + monotonic() + + settings.query_metrics_snapshot_poll_seconds + ) + except Exception: + retry_not_before[retry_key] = ( + monotonic() + + settings.query_metrics_snapshot_retry_seconds + ) + logger.exception( + "snapshot refresh failed family=%s window=%s", + family, + window, + ) + if not attempted: + stop.wait(settings.query_metrics_snapshot_poll_seconds) + except Exception: + logger.exception("snapshot worker connection loop failed") + stop.wait(settings.query_metrics_snapshot_retry_seconds) + + +def main() -> None: + logging.basicConfig( + level=getattr(logging, get_settings().log_level.upper(), logging.INFO), + format="%(asctime)s %(levelname)s %(name)s %(message)s", + ) + settings = get_settings() + stop = Event() + + def request_stop(_: int, __: object) -> None: + stop.set() + + signal.signal(signal.SIGTERM, request_stop) + signal.signal(signal.SIGINT, request_stop) + run_worker(settings, stop) + + +if __name__ == "__main__": + main() diff --git a/backend/app/version.py b/backend/app/version.py index f896033..735615f 100644 --- a/backend/app/version.py +++ b/backend/app/version.py @@ -1,4 +1,4 @@ """Release identifiers shared by every backend process.""" -APPLICATION_VERSION = "1.1.0" -EXPECTED_MIGRATION = "0014" +APPLICATION_VERSION = "1.1.1" +EXPECTED_MIGRATION = "0016" diff --git a/backend/tests/test_global_trend_cache.py b/backend/tests/test_global_trend_cache.py index 774ce42..de2c15e 100644 --- a/backend/tests/test_global_trend_cache.py +++ b/backend/tests/test_global_trend_cache.py @@ -66,6 +66,67 @@ async def wait_until(predicate: Any) -> None: raise AssertionError("asynchronous cache condition was not reached") +@pytest.mark.asyncio +async def test_global_trend_loader_reads_persistent_scoped_snapshot( + monkeypatch: pytest.MonkeyPatch, +) -> None: + class Cursor: + def __init__(self) -> None: + self.executions: list[tuple[str, list[Any] | tuple[Any, ...]]] = [] + + async def __aenter__(self) -> Cursor: + return self + + async def __aexit__(self, *_: object) -> None: + return None + + async def execute( + self, query: str, params: list[Any] | tuple[Any, ...] + ) -> None: + self.executions.append((query, params)) + + async def fetchall(self) -> list[dict[str, Any]]: + return [trend_row(1)] + + class Connection: + def __init__(self) -> None: + self.control = Cursor() + self.data = Cursor() + self.cursor_count = 0 + + async def __aenter__(self) -> Connection: + return self + + async def __aexit__(self, *_: object) -> None: + return None + + def cursor(self) -> Cursor: + self.cursor_count += 1 + return self.control if self.cursor_count == 1 else self.data + + class Pool: + def __init__(self) -> None: + self.connection_instance = Connection() + + def connection(self) -> Connection: + return self.connection_instance + + fake_pool = Pool() + monkeypatch.setattr(powa_module, "pool", fake_pool) + repository = cache_repository(FakeClock()) + + rows = await repository._load_global_trend_snapshot( + "24h", server_id=7, database_id=16_384 + ) + + query, params = fake_pool.connection_instance.data.executions[0] + assert rows == [trend_row(1)] + assert "advisor.global_trend_snapshot_24h" in query + assert "advisor.query_trend(" not in query + assert "IS NOT DISTINCT FROM" in query + assert params == (7, 16_384, 100_001) + + @pytest.mark.asyncio async def test_global_trend_cold_singleflight_is_shielded_and_deep_copy_safe( monkeypatch: pytest.MonkeyPatch, @@ -220,7 +281,7 @@ async def load(_: str) -> list[dict[str, Any]]: @pytest.mark.asyncio -async def test_query_metrics_and_global_trend_share_repository_refresh_lock( +async def test_precomputed_metrics_and_trend_snapshots_read_concurrently( monkeypatch: pytest.MonkeyPatch, ) -> None: repository = cache_repository(FakeClock()) @@ -243,8 +304,8 @@ async def load_trend(_: str) -> list[dict[str, Any]]: metrics_request = asyncio.create_task(repository.query_rows(window="24h")) await metrics_started.wait() trend_request = asyncio.create_task(repository.trend(window="24h")) - await asyncio.sleep(0) - assert not trend_started.is_set() + await wait_until(trend_started.is_set) + assert trend_started.is_set() assert repository._query_metrics_refresh_lock is repository._repository_refresh_lock metrics_release.set() diff --git a/backend/tests/test_product_api.py b/backend/tests/test_product_api.py index 0e29aac..c031aaa 100644 --- a/backend/tests/test_product_api.py +++ b/backend/tests/test_product_api.py @@ -488,17 +488,17 @@ async def rows(**_: object): def test_release_contract_reports_expected_migration() -> None: payload = _release_payload({ - "current_migration": "0014", - "applied_count": 14, - "latest_applied_at": datetime(2026, 7, 26, tzinfo=timezone.utc), + "current_migration": "0016", + "applied_count": 16, + "latest_applied_at": datetime(2026, 7, 30, tzinfo=timezone.utc), }) - assert payload["applicationVersion"] == "1.1.0" - assert payload["migration"]["expected"] == "0014" + assert payload["applicationVersion"] == "1.1.1" + assert payload["migration"]["expected"] == "0016" assert payload["migration"]["upToDate"] is True def test_backend_versions_and_scoped_openapi_contract_are_aligned() -> None: - assert api_app.version == evaluator_app.version == clone_app.version == "1.1.0" + assert api_app.version == evaluator_app.version == clone_app.version == "1.1.1" schema = api_app.openapi() for path in ("/api/v1/system-health", "/api/v1/operations"): parameters = { diff --git a/backend/tests/test_query_rows_cache.py b/backend/tests/test_query_rows_cache.py index 4de00af..befb7bb 100644 --- a/backend/tests/test_query_rows_cache.py +++ b/backend/tests/test_query_rows_cache.py @@ -303,7 +303,7 @@ async def load(_: str) -> list[dict[str, Any]]: @pytest.mark.asyncio -async def test_different_windows_do_not_scan_repository_concurrently( +async def test_different_windows_read_precomputed_snapshots_concurrently( monkeypatch: pytest.MonkeyPatch, ) -> None: repository = cache_repository(FakeClock()) @@ -322,8 +322,8 @@ async def load(window: str) -> list[dict[str, Any]]: first = asyncio.create_task(repository.query_rows(window="1h")) await first_started.wait() second = asyncio.create_task(repository.query_rows(window="24h")) - await asyncio.sleep(0) - assert started_windows == ["1h"] + await wait_until(lambda: started_windows == ["1h", "24h"]) + assert started_windows == ["1h", "24h"] first_release.set() await asyncio.gather(first, second) @@ -641,12 +641,15 @@ def connection(self) -> Connection: with pytest.raises(QueryMetricsSnapshotTooLarge): await repository._load_query_metrics_snapshot("1h") - assert connection.cursor_names == [None, "advisor_query_metrics_cache_refresh"] + assert connection.cursor_names == [None, "advisor_query_metrics_snapshot_read"] assert "SET LOCAL application_name" in connection.control_cursor.query - assert "advisor-query-metrics-cache-refresh" in cursor.query + assert "advisor-query-metrics-snapshot-read" in cursor.query + assert "advisor.query_metrics_snapshot_1h" in cursor.query + assert "advisor.query_metrics(" not in cursor.query + assert "advisor.query_annotations" in cursor.query assert "pg_column_size(metrics)" in cursor.query assert "LIMIT %s" in cursor.query - assert cursor.params == ("1 hour", 3) + assert cursor.params == (3,) assert cursor.fetch_sizes == [3, 1] diff --git a/backend/tests/test_snapshot_worker.py b/backend/tests/test_snapshot_worker.py new file mode 100644 index 0000000..ba987c9 --- /dev/null +++ b/backend/tests/test_snapshot_worker.py @@ -0,0 +1,78 @@ +from __future__ import annotations + +from datetime import datetime, timedelta, timezone +from typing import Any + +from app.config import ( + GLOBAL_TREND_SNAPSHOT_VIEWS, + QUERY_METRICS_SNAPSHOT_VIEWS, + Settings, + WINDOW_INTERVALS, +) +from app.snapshot_worker import ( + WINDOW_ORDER, + due_windows, + global_trend_due_windows, + refresh_intervals, +) + + +class Result: + def __init__(self, rows: list[dict[str, Any]]) -> None: + self.rows = rows + + def fetchall(self) -> list[dict[str, Any]]: + return self.rows + + +class Connection: + def __init__(self, rows: list[dict[str, Any]]) -> None: + self.rows = rows + + def execute(self, _: str) -> Result: + return Result(self.rows) + + +def test_snapshot_view_mapping_covers_every_supported_window() -> None: + assert tuple(QUERY_METRICS_SNAPSHOT_VIEWS) == WINDOW_ORDER + assert tuple(GLOBAL_TREND_SNAPSHOT_VIEWS) == WINDOW_ORDER + assert set(QUERY_METRICS_SNAPSHOT_VIEWS) == set(WINDOW_INTERVALS) + assert set(GLOBAL_TREND_SNAPSHOT_VIEWS) == set(WINDOW_INTERVALS) + + +def test_due_windows_prioritize_short_windows_and_respect_cadence() -> None: + observed_at = datetime(2026, 7, 30, 12, tzinfo=timezone.utc) + settings = Settings( + query_metrics_snapshot_1h_refresh_seconds=15 * 60, + query_metrics_snapshot_24h_refresh_seconds=60 * 60, + query_metrics_snapshot_7d_refresh_seconds=6 * 60 * 60, + query_metrics_snapshot_30d_refresh_seconds=12 * 60 * 60, + ) + connection = Connection( + [ + { + "window_key": "1h", + "refreshed_at": observed_at - timedelta(minutes=20), + }, + { + "window_key": "24h", + "refreshed_at": observed_at - timedelta(minutes=30), + }, + {"window_key": "7d", "refreshed_at": None}, + { + "window_key": "30d", + "refreshed_at": observed_at - timedelta(hours=1), + }, + ] + ) + + assert due_windows(connection, settings, now=observed_at) == ["1h", "7d"] + assert global_trend_due_windows( + connection, settings, now=observed_at + ) == ["1h", "7d"] + assert refresh_intervals(settings) == { + "1h": 900, + "24h": 3600, + "7d": 21600, + "30d": 43200, + } diff --git a/backend/tests/test_trend_repository.py b/backend/tests/test_trend_repository.py index 585dd0a..e7bf664 100644 --- a/backend/tests/test_trend_repository.py +++ b/backend/tests/test_trend_repository.py @@ -54,7 +54,7 @@ def connection(self) -> FakeConnection: @pytest.mark.asyncio -async def test_global_trend_uses_preaggregated_two_argument_helper( +async def test_global_trend_reads_precomputed_fleet_snapshot( monkeypatch: pytest.MonkeyPatch, ) -> None: cursor = FakeCursor() @@ -66,14 +66,13 @@ async def test_global_trend_uses_preaggregated_two_argument_helper( assert rows[0]["calls"] == 3 tag_query, tag_params = cursor.executions[0] assert "SET LOCAL application_name" in tag_query - assert "advisor-global-trend-cache-refresh" in tag_query + assert "advisor-global-trend-snapshot-read" in tag_query assert tag_params == [] query, params = cursor.executions[1] - assert "advisor.query_trend(" in query - assert "now() - %s::interval" in query - assert "advisor.query_deltas" not in query - assert "advisor-global-trend-cache-refresh" in query - assert params == ["1 hour", "5 minutes", 100_001] + assert "advisor.global_trend_snapshot_1h" in query + assert "advisor.query_trend(" not in query + assert "advisor-global-trend-snapshot-read" in query + assert params == [None, None, 100_001] @pytest.mark.asyncio @@ -100,7 +99,7 @@ async def test_query_trend_pushes_complete_scope_into_five_argument_helper( @pytest.mark.asyncio -async def test_server_trend_uses_four_argument_scope_helper( +async def test_server_trend_reads_precomputed_server_scope( monkeypatch: pytest.MonkeyPatch, ) -> None: cursor = FakeCursor() @@ -110,11 +109,12 @@ async def test_server_trend_uses_four_argument_scope_helper( await repository.trend(window="1h", server_id=7) tag_query, tag_params = cursor.executions[0] - assert "advisor-global-trend-cache-refresh" in tag_query + assert "advisor-global-trend-snapshot-read" in tag_query assert tag_params == [] query, params = cursor.executions[1] - assert "now() - %s::interval, %s::interval, %s, %s)" in query - assert params == ["1 hour", "5 minutes", 7, None, 100_001] + assert "advisor.global_trend_snapshot_1h" in query + assert "advisor.query_trend(" not in query + assert params == [7, None, 100_001] @pytest.mark.asyncio diff --git a/compose.production.yaml b/compose.production.yaml new file mode 100644 index 0000000..f12bf92 --- /dev/null +++ b/compose.production.yaml @@ -0,0 +1,36 @@ +services: + workload: + restart: unless-stopped + + caddy: + image: caddy:2.10.2-alpine + environment: + DASHBOARD_BASIC_AUTH_USER: ${DASHBOARD_BASIC_AUTH_USER:?Set DASHBOARD_BASIC_AUTH_USER in .env} + DASHBOARD_BASIC_AUTH_HASH: ${DASHBOARD_BASIC_AUTH_HASH:?Set DASHBOARD_BASIC_AUTH_HASH in .env} + ports: + - "80:80" + - "443:443" + - "443:443/udp" + volumes: + - ./deployment/production/Caddyfile:/etc/caddy/Caddyfile:ro + - caddy_data:/data + - caddy_config:/config + depends_on: + web: + condition: service_healthy + restart: unless-stopped + read_only: true + tmpfs: + - /tmp:rw,noexec,nosuid,size=32m + - /run:rw,noexec,nosuid,size=8m + cap_drop: [ALL] + cap_add: [NET_BIND_SERVICE] + security_opt: + - no-new-privileges:true + pids_limit: 128 + mem_limit: 256m + networks: [advisor] + +volumes: + caddy_data: + caddy_config: diff --git a/compose.yaml b/compose.yaml index f726d74..914778c 100644 --- a/compose.yaml +++ b/compose.yaml @@ -138,6 +138,44 @@ services: restart: "no" networks: [advisor] + query-metrics-snapshot-worker: + build: + context: ./backend + image: postgresql-advisor/api:iteration-2.7 + command: ["python", "-m", "app.snapshot_worker"] + environment: + DATABASE_URL: ${ADVISOR_API_DATABASE_URL:-} + DATABASE_HOST: ${ADVISOR_API_DATABASE_HOST:-repository-db} + DATABASE_PORT: ${ADVISOR_API_DATABASE_PORT:-5433} + DATABASE_NAME: ${ADVISOR_API_DATABASE_NAME:-powa_repository} + DATABASE_USER: ${ADVISOR_API_DATABASE_USER:-advisor_api} + DATABASE_PASSWORD: ${ADVISOR_API_PASSWORD:?Set ADVISOR_API_PASSWORD in .env} + DATABASE_SSLMODE: ${ADVISOR_API_DATABASE_SSLMODE:-disable} + QUERY_METRICS_SNAPSHOT_POLL_SECONDS: ${QUERY_METRICS_SNAPSHOT_POLL_SECONDS:-15} + QUERY_METRICS_SNAPSHOT_1H_REFRESH_SECONDS: ${QUERY_METRICS_SNAPSHOT_1H_REFRESH_SECONDS:-900} + QUERY_METRICS_SNAPSHOT_24H_REFRESH_SECONDS: ${QUERY_METRICS_SNAPSHOT_24H_REFRESH_SECONDS:-3600} + QUERY_METRICS_SNAPSHOT_7D_REFRESH_SECONDS: ${QUERY_METRICS_SNAPSHOT_7D_REFRESH_SECONDS:-21600} + QUERY_METRICS_SNAPSHOT_30D_REFRESH_SECONDS: ${QUERY_METRICS_SNAPSHOT_30D_REFRESH_SECONDS:-43200} + QUERY_METRICS_SNAPSHOT_STATEMENT_TIMEOUT_SECONDS: ${QUERY_METRICS_SNAPSHOT_STATEMENT_TIMEOUT_SECONDS:-1800} + QUERY_METRICS_SNAPSHOT_RETRY_SECONDS: ${QUERY_METRICS_SNAPSHOT_RETRY_SECONDS:-60} + LOG_LEVEL: ${LOG_LEVEL:-INFO} + depends_on: + repository-db: + condition: service_healthy + repository-migrate: + condition: service_completed_successfully + init: true + read_only: true + tmpfs: + - /tmp:rw,noexec,nosuid,size=16m + cap_drop: [ALL] + security_opt: + - no-new-privileges:true + pids_limit: 32 + mem_limit: ${QUERY_METRICS_SNAPSHOT_WORKER_MEMORY_LIMIT:-256m} + restart: unless-stopped + networks: [advisor] + collector: build: context: . @@ -166,6 +204,9 @@ services: condition: service_healthy repository-migrate: condition: service_completed_successfully + read_only: true + tmpfs: + - /tmp:rw,noexec,nosuid,size=16m,mode=1777 healthcheck: test: ["CMD-SHELL", "python -c \"import os,signal; os.kill(1, signal.SIGCONT)\""] interval: 15s @@ -275,6 +316,8 @@ services: condition: service_healthy repository-migrate: condition: service_completed_successfully + query-metrics-snapshot-worker: + condition: service_started collector: condition: service_started evaluator: @@ -448,6 +491,10 @@ services: WORKLOAD_PROFILE: ${WORKLOAD_PROFILE:-normal} WORKLOAD_DURATION_SECONDS: ${WORKLOAD_DURATION_SECONDS:-0} WORKLOAD_INTERVAL_SECONDS: ${WORKLOAD_INTERVAL_SECONDS:-0.25} + WORKLOAD_INTERVAL_JITTER_RATIO: ${WORKLOAD_INTERVAL_JITTER_RATIO:-0} + WORKLOAD_TRAFFIC_PHASE_SECONDS: ${WORKLOAD_TRAFFIC_PHASE_SECONDS:-0} + WORKLOAD_TRAFFIC_MIN_INTERVAL_MULTIPLIER: ${WORKLOAD_TRAFFIC_MIN_INTERVAL_MULTIPLIER:-1} + WORKLOAD_TRAFFIC_MAX_INTERVAL_MULTIPLIER: ${WORKLOAD_TRAFFIC_MAX_INTERVAL_MULTIPLIER:-1} WORKLOAD_WORKERS: ${WORKLOAD_WORKERS:-6} WORKLOAD_REPORT_INTERVAL_SECONDS: ${WORKLOAD_REPORT_INTERVAL_SECONDS:-10} WORKLOAD_RANDOM_SEED: ${WORKLOAD_RANDOM_SEED:-20260725} diff --git a/deployment/production/Caddyfile b/deployment/production/Caddyfile new file mode 100644 index 0000000..956b1f6 --- /dev/null +++ b/deployment/production/Caddyfile @@ -0,0 +1,23 @@ +{ + email ops@trades.engineer + admin off +} + +trades.engineer, trades.164-92-171-43.sslip.io { + encode zstd gzip + + basic_auth { + {$DASHBOARD_BASIC_AUTH_USER} {$DASHBOARD_BASIC_AUTH_HASH} + } + + header { + Strict-Transport-Security "max-age=31536000; includeSubDomains" + X-Content-Type-Options "nosniff" + X-Frame-Options "DENY" + Referrer-Policy "strict-origin-when-cross-origin" + Permissions-Policy "camera=(), geolocation=(), microphone=()" + -Server + } + + reverse_proxy web:80 +} diff --git a/deployment/production/bootstrap-env.sh b/deployment/production/bootstrap-env.sh new file mode 100755 index 0000000..0e54417 --- /dev/null +++ b/deployment/production/bootstrap-env.sh @@ -0,0 +1,111 @@ +#!/usr/bin/env bash +set -Eeuo pipefail + +project_dir="$(cd -- "$(dirname -- "${BASH_SOURCE[0]}")/../.." && pwd -P)" +cd "$project_dir" + +if [[ -e .env ]]; then + printf 'Refusing to overwrite existing %s/.env\n' "$project_dir" >&2 + exit 1 +fi + +command -v docker >/dev/null 2>&1 || { + printf 'docker is required\n' >&2 + exit 1 +} +command -v openssl >/dev/null 2>&1 || { + printf 'openssl is required\n' >&2 + exit 1 +} + +umask 077 + +random_hex() { + openssl rand -hex 32 +} + +basic_auth_user="trades" +basic_auth_password="$(openssl rand -base64 24 | tr -d '\n' | tr '/+' '_-')" +admin_token="$(openssl rand -base64 32 | tr -d '\n' | tr '/+' '_-')" +admin_token_sha256="$(printf '%s' "$admin_token" | openssl dgst -sha256 -r | awk '{print $1}')" +basic_auth_hash="$( + printf '%s\n' "$basic_auth_password" | + docker run --rm -i caddy:2.10.2-alpine \ + caddy hash-password +)" +principal_json="$( + printf '[{"credential_id":"prod-admin","subject":"trades-engineer-admin","token_sha256":"%s","roles":["analyst","annotator","admin"]}]' \ + "$admin_token_sha256" +)" + +{ + printf 'COMPOSE_PROJECT_NAME=trades-engineer\n' + printf 'COMPOSE_FILE=compose.yaml:compose.production.yaml\n' + printf 'POSTGRES_ADMIN_PASSWORD=%s\n' "$(random_hex)" + printf 'POWA_COLLECTOR_PASSWORD=%s\n' "$(random_hex)" + printf 'ADVISOR_API_PASSWORD=%s\n' "$(random_hex)" + printf 'ADVISOR_EVALUATOR_PASSWORD=%s\n' "$(random_hex)" + printf 'ADVISOR_EVALUATOR_READ_SCHEMAS=public\n' + printf 'EVALUATOR_TOKEN=%s\n' "$(random_hex)" + printf 'ADVISOR_JOIN_SOURCE_PASSWORD=%s\n' "$(random_hex)" + printf 'ADVISOR_JOIN_REPOSITORY_PASSWORD=%s\n' "$(random_hex)" + printf 'WORKLOAD_DB_PASSWORD=%s\n' "$(random_hex)" + printf 'CLONE_ADMIN_PASSWORD=%s\n' "$(random_hex)" + printf 'CLONE_RUNNER_PASSWORD=%s\n' "$(random_hex)" + printf 'CLONE_EVALUATOR_TOKEN=%s\n' "$(random_hex)" + printf "ADVISOR_AUTH_PRINCIPALS='%s'\n" "$principal_json" + printf 'SOURCE_DB_BIND=127.0.0.1\n' + printf 'SOURCE_DB_PORT=15432\n' + printf 'REPOSITORY_DB_BIND=127.0.0.1\n' + printf 'REPOSITORY_DB_PORT=15433\n' + printf 'API_BIND=127.0.0.1\n' + printf 'API_PORT=8000\n' + printf 'WEB_BIND=127.0.0.1\n' + printf 'WEB_PORT=5173\n' + printf 'DEFAULT_WINDOW=24h\n' + printf 'RETENTION_DAYS=90\n' + printf 'LOG_LEVEL=INFO\n' + printf 'API_MEMORY_LIMIT=1g\n' + printf 'QUERY_METRICS_SNAPSHOT_POLL_SECONDS=15\n' + printf 'QUERY_METRICS_SNAPSHOT_1H_REFRESH_SECONDS=900\n' + printf 'QUERY_METRICS_SNAPSHOT_24H_REFRESH_SECONDS=3600\n' + printf 'QUERY_METRICS_SNAPSHOT_7D_REFRESH_SECONDS=21600\n' + printf 'QUERY_METRICS_SNAPSHOT_30D_REFRESH_SECONDS=43200\n' + printf 'QUERY_METRICS_SNAPSHOT_STATEMENT_TIMEOUT_SECONDS=1800\n' + printf 'QUERY_METRICS_SNAPSHOT_RETRY_SECONDS=60\n' + printf 'QUERY_METRICS_SNAPSHOT_WORKER_MEMORY_LIMIT=256m\n' + printf 'SOURCE_DB_SHM_SIZE=512mb\n' + printf 'REPOSITORY_DB_SHM_SIZE=256mb\n' + printf 'WORKLOAD_PROFILE=erp\n' + printf 'WORKLOAD_DURATION_SECONDS=7200\n' + printf 'WORKLOAD_WORKERS=3\n' + printf 'WORKLOAD_INTERVAL_SECONDS=0.20\n' + printf 'WORKLOAD_INTERVAL_JITTER_RATIO=0.30\n' + printf 'WORKLOAD_TRAFFIC_PHASE_SECONDS=45\n' + printf 'WORKLOAD_TRAFFIC_MIN_INTERVAL_MULTIPLIER=0.90\n' + printf 'WORKLOAD_TRAFFIC_MAX_INTERVAL_MULTIPLIER=3.00\n' + printf 'WORKLOAD_STATEMENT_TIMEOUT_MS=30000\n' + printf 'WORKLOAD_LOCK_TIMEOUT_MS=1000\n' + printf 'WORKLOAD_LOCK_HOLD_MS=30\n' + printf 'WORKLOAD_ERP_TABLE_COUNT=500\n' + printf 'WORKLOAD_ERP_QUERY_VARIANTS_PER_TABLE=8\n' + printf 'WORKLOAD_ERP_ROWS_PER_TABLE=2000\n' + printf 'REGISTER_DEMO_SOURCE=true\n' + printf 'DEMO_SOURCE_FREQUENCY=60\n' + printf 'POWA_SOURCE_SSLMODE=prefer\n' + printf 'DASHBOARD_BASIC_AUTH_USER=%s\n' "$basic_auth_user" + printf "DASHBOARD_BASIC_AUTH_HASH='%s'\n" "$basic_auth_hash" +} > .env +chmod 0600 .env + +credentials_path="${1:-/home/deploy/.trades-engineer-credentials}" +{ + printf 'Dashboard URL: https://trades.engineer\n' + printf 'Dashboard basic-auth user: %s\n' "$basic_auth_user" + printf 'Dashboard basic-auth password: %s\n' "$basic_auth_password" + printf 'Advisor API admin bearer token: %s\n' "$admin_token" +} > "$credentials_path" +chmod 0600 "$credentials_path" + +docker compose --profile realistic-load config --quiet +printf 'Created %s/.env and %s\n' "$project_dir" "$credentials_path" diff --git a/deployment/workload/test_workload.py b/deployment/workload/test_workload.py index 48aff82..faf43a8 100644 --- a/deployment/workload/test_workload.py +++ b/deployment/workload/test_workload.py @@ -137,6 +137,10 @@ def test_documented_environment_overrides_are_applied(self) -> None: WORKLOAD_DURATION_SECONDS="60", WORKLOAD_WORKERS="12", WORKLOAD_INTERVAL_SECONDS="0.125", + WORKLOAD_INTERVAL_JITTER_RATIO="0.4", + WORKLOAD_TRAFFIC_PHASE_SECONDS="30", + WORKLOAD_TRAFFIC_MIN_INTERVAL_MULTIPLIER="0.5", + WORKLOAD_TRAFFIC_MAX_INTERVAL_MULTIPLIER="2.5", WORKLOAD_REPORT_INTERVAL_SECONDS="4", WORKLOAD_RANDOM_SEED="42", WORKLOAD_STATEMENT_TIMEOUT_MS="7000", @@ -146,6 +150,10 @@ def test_documented_environment_overrides_are_applied(self) -> None: self.assertEqual(config.duration_seconds, 60) self.assertEqual(config.workers, 12) self.assertEqual(config.interval_seconds, 0.125) + self.assertEqual(config.interval_jitter_ratio, 0.4) + self.assertEqual(config.traffic_phase_seconds, 30) + self.assertEqual(config.traffic_min_interval_multiplier, 0.5) + self.assertEqual(config.traffic_max_interval_multiplier, 2.5) self.assertEqual(config.report_interval_seconds, 4) self.assertEqual(config.random_seed, 42) self.assertEqual(config.statement_timeout_ms, 7000) @@ -159,6 +167,12 @@ def test_invalid_configuration_fails_closed(self) -> None: {"WORKLOAD_WORKERS": "2"}, {"WORKLOAD_DURATION_SECONDS": "-1"}, {"WORKLOAD_INTERVAL_SECONDS": "nan"}, + {"WORKLOAD_INTERVAL_JITTER_RATIO": "1"}, + {"WORKLOAD_TRAFFIC_PHASE_SECONDS": "-1"}, + { + "WORKLOAD_TRAFFIC_MIN_INTERVAL_MULTIPLIER": "3", + "WORKLOAD_TRAFFIC_MAX_INTERVAL_MULTIPLIER": "2", + }, {"WORKLOAD_ERP_TABLE_COUNT": "501"}, {"WORKLOAD_ERP_QUERY_VARIANTS_PER_TABLE": "9"}, {"WORKLOAD_ERP_ROWS_PER_TABLE": "5001"}, @@ -329,6 +343,31 @@ def test_weighted_operation_selection_never_crosses_role(self) -> None: } self.assertEqual(selected, {role}) + def test_variable_traffic_interval_is_bounded_and_phase_based(self) -> None: + config = test_config( + WORKLOAD_INTERVAL_SECONDS="0.2", + WORKLOAD_INTERVAL_JITTER_RATIO="0.25", + WORKLOAD_TRAFFIC_PHASE_SECONDS="30", + WORKLOAD_TRAFFIC_MIN_INTERVAL_MULTIPLIER="0.5", + WORKLOAD_TRAFFIC_MAX_INTERVAL_MULTIPLIER="2", + ) + first_phase = workload.traffic_phase_interval_multiplier(config, 5) + self.assertEqual( + first_phase, + workload.traffic_phase_interval_multiplier(config, 29.9), + ) + self.assertNotEqual( + first_phase, + workload.traffic_phase_interval_multiplier(config, 30), + ) + + intervals = [ + workload.traffic_interval_seconds(config, random.Random(seed), 5) + for seed in range(20) + ] + self.assertGreater(len(set(intervals)), 1) + self.assertTrue(all(0.075 <= interval <= 0.5 for interval in intervals)) + class OperationTests(unittest.TestCase): def test_mutation_write_uses_bounded_cleanup(self) -> None: @@ -528,7 +567,7 @@ def test_timeout_operational_errors_are_not_connection_failures(self) -> None: workload._is_connection_failure(connection_error, connection) ) - def test_start_and_final_use_utc_timestamps_without_changing_heartbeat(self) -> None: + def test_start_final_and_heartbeat_have_stable_metadata(self) -> None: emitted: list[dict[str, object]] = [] config = test_config() @@ -574,10 +613,20 @@ def test_start_and_final_use_utc_timestamps_without_changing_heartbeat(self) -> "profile", "elapsedSeconds", "remainingSeconds", + "traffic", "totals", "categories", }, ) + self.assertEqual( + heartbeat["traffic"], + { + "baseIntervalSeconds": config.interval_seconds, + "intervalJitterRatio": config.interval_jitter_ratio, + "phaseSeconds": config.traffic_phase_seconds, + "phaseIntervalMultiplier": 1.0, + }, + ) failed_emissions: list[dict[str, object]] = [] with ( diff --git a/deployment/workload/workload.py b/deployment/workload/workload.py index 3dccde7..696cc92 100644 --- a/deployment/workload/workload.py +++ b/deployment/workload/workload.py @@ -558,6 +558,10 @@ class WorkloadConfig: duration_seconds: int workers: int interval_seconds: float + interval_jitter_ratio: float + traffic_phase_seconds: int + traffic_min_interval_multiplier: float + traffic_max_interval_multiplier: float report_interval_seconds: int random_seed: int statement_timeout_ms: int @@ -598,6 +602,39 @@ def from_env(cls, environ: Mapping[str, str] | None = None) -> WorkloadConfig: interval_seconds = _number( values, defaults, "WORKLOAD_INTERVAL_SECONDS", 0, 60 ) + interval_jitter_ratio = _number( + values, + {"WORKLOAD_INTERVAL_JITTER_RATIO": "0"}, + "WORKLOAD_INTERVAL_JITTER_RATIO", + 0, + 0.95, + ) + traffic_phase_seconds = _integer( + values, + {"WORKLOAD_TRAFFIC_PHASE_SECONDS": "0"}, + "WORKLOAD_TRAFFIC_PHASE_SECONDS", + 0, + 3_600, + ) + traffic_min_interval_multiplier = _number( + values, + {"WORKLOAD_TRAFFIC_MIN_INTERVAL_MULTIPLIER": "1"}, + "WORKLOAD_TRAFFIC_MIN_INTERVAL_MULTIPLIER", + 0.1, + 5, + ) + traffic_max_interval_multiplier = _number( + values, + {"WORKLOAD_TRAFFIC_MAX_INTERVAL_MULTIPLIER": "1"}, + "WORKLOAD_TRAFFIC_MAX_INTERVAL_MULTIPLIER", + 0.1, + 5, + ) + if traffic_min_interval_multiplier > traffic_max_interval_multiplier: + raise ValueError( + "WORKLOAD_TRAFFIC_MIN_INTERVAL_MULTIPLIER cannot exceed " + "WORKLOAD_TRAFFIC_MAX_INTERVAL_MULTIPLIER" + ) report_interval_seconds = _integer( values, defaults, "WORKLOAD_REPORT_INTERVAL_SECONDS", 1, 300 ) @@ -703,6 +740,10 @@ def from_env(cls, environ: Mapping[str, str] | None = None) -> WorkloadConfig: duration_seconds=duration_seconds, workers=workers, interval_seconds=interval_seconds, + interval_jitter_ratio=interval_jitter_ratio, + traffic_phase_seconds=traffic_phase_seconds, + traffic_min_interval_multiplier=traffic_min_interval_multiplier, + traffic_max_interval_multiplier=traffic_max_interval_multiplier, report_interval_seconds=report_interval_seconds, random_seed=random_seed, statement_timeout_ms=statement_timeout_ms, @@ -1332,6 +1373,34 @@ def worker_seed(base_seed: int, worker_id: int) -> int: return (base_seed + worker_id * 1_000_003) % (2**63) +def traffic_phase_interval_multiplier( + config: WorkloadConfig, elapsed_seconds: float +) -> float: + if config.traffic_phase_seconds == 0: + return 1.0 + phase_index = max(0, int(elapsed_seconds // config.traffic_phase_seconds)) + phase_seed = ( + config.random_seed ^ ((phase_index + 1) * 0x9E37_79B9_7F4A_7C15) + ) % (2**63) + return random.Random(phase_seed).uniform( + config.traffic_min_interval_multiplier, + config.traffic_max_interval_multiplier, + ) + + +def traffic_interval_seconds( + config: WorkloadConfig, rng: random.Random, elapsed_seconds: float +) -> float: + if config.interval_seconds == 0: + return 0.0 + jitter_multiplier = rng.uniform( + 1 - config.interval_jitter_ratio, + 1 + config.interval_jitter_ratio, + ) + phase_multiplier = traffic_phase_interval_multiplier(config, elapsed_seconds) + return min(60.0, config.interval_seconds * jitter_multiplier * phase_multiplier) + + ROLE_SHARES: dict[str, dict[str, float]] = { "quick": {READER_ROLE: 0.50, REPORTER_ROLE: 0.30, WRITER_ROLE: 0.20}, "normal": {READER_ROLE: 0.50, REPORTER_ROLE: 0.30, WRITER_ROLE: 0.20}, @@ -1699,6 +1768,7 @@ def run_worker( deadline: float, ) -> None: rng = random.Random(worker_seed(config.random_seed, worker_id)) + traffic_started_at = time.monotonic() reconnect_delay = 0.25 while not stop_event.is_set() and time.monotonic() < deadline: try: @@ -1774,8 +1844,13 @@ def run_worker( metrics.record_erp_fingerprint(erp_ordinal) remaining = deadline - time.monotonic() - if config.interval_seconds > 0 and remaining > 0: - stop_event.wait(min(config.interval_seconds, remaining)) + interval_seconds = traffic_interval_seconds( + config, + rng, + max(0.0, time.monotonic() - traffic_started_at), + ) + if interval_seconds > 0 and remaining > 0: + stop_event.wait(min(interval_seconds, remaining)) except (psycopg.OperationalError, OSError) as exc: metrics.record_connection_error(getattr(exc, "sqlstate", None)) remaining = deadline - time.monotonic() @@ -1817,6 +1892,14 @@ def _heartbeat_payload( if math.isinf(remaining_seconds) else round(max(0, remaining_seconds), 3) ), + "traffic": { + "baseIntervalSeconds": config.interval_seconds, + "intervalJitterRatio": config.interval_jitter_ratio, + "phaseSeconds": config.traffic_phase_seconds, + "phaseIntervalMultiplier": round( + traffic_phase_interval_multiplier(config, elapsed_seconds), 3 + ), + }, "totals": snapshot["totals"], "categories": snapshot["categories"], } @@ -1862,6 +1945,15 @@ def request_stop(signum: int, _frame: object) -> None: "workers": config.workers, "roleWorkers": role_workers, "randomSeed": config.random_seed, + "intervalSeconds": config.interval_seconds, + "intervalJitterRatio": config.interval_jitter_ratio, + "trafficPhaseSeconds": config.traffic_phase_seconds, + "trafficMinIntervalMultiplier": ( + config.traffic_min_interval_multiplier + ), + "trafficMaxIntervalMultiplier": ( + config.traffic_max_interval_multiplier + ), "dataBounds": bounds.__dict__, "sqlTemplateCount": len(SQL_TEMPLATES), "erpTableCount": config.erp_table_count, @@ -1962,6 +2054,15 @@ def request_stop(signum: int, _frame: object) -> None: "workers": config.workers, "roleWorkers": role_workers, "randomSeed": config.random_seed, + "intervalSeconds": config.interval_seconds, + "intervalJitterRatio": config.interval_jitter_ratio, + "trafficPhaseSeconds": config.traffic_phase_seconds, + "trafficMinIntervalMultiplier": ( + config.traffic_min_interval_multiplier + ), + "trafficMaxIntervalMultiplier": ( + config.traffic_max_interval_multiplier + ), "dataBounds": bounds.__dict__, "sqlTemplateCount": len(SQL_TEMPLATES), "erpTableCount": config.erp_table_count, diff --git a/docs/REPOSITORY_MIGRATIONS.md b/docs/REPOSITORY_MIGRATIONS.md index 818fd56..b2a9af8 100644 --- a/docs/REPOSITORY_MIGRATIONS.md +++ b/docs/REPOSITORY_MIGRATIONS.md @@ -87,6 +87,17 @@ release bilgisini ve tablo yazma maliyeti sinyalini ekler. Bu migration uygulama `1.1.0` ile birlikte yükseltilir; geri dönüş sınırı [upgrade/rollback runbook'undadır](UPGRADE_ROLLBACK.md). +`0015`, dört sabit dashboard penceresi için kalıcı query-metrics materialized +snapshot'ları ve yenileme durum tablosunu ekler. İlk population worker tarafından +migration dışında yapılır. Sonraki yenilemeler `CONCURRENTLY` çalıştığından +hesaplama tamamlanana kadar önceki atomik sonuç okunur; API request yolu +`advisor.query_metrics(interval)` çağırmaz. + +`0016`, overview trendlerini database scope'ta bir kez hesaplayan ve aynı bucket +satırlarından server/global toplamlarını türeten dört kalıcı materialized +snapshot ekler. Böylece global, server ve database overview yolları cold API +başlangıcında dahi `advisor.query_trend(...)` çalıştırmaz. + ## Yeni migration ekleme 1. Mevcut migration dosyalarını değiştirmeyin. Özellikle diff --git a/docs/UPGRADE_ROLLBACK.md b/docs/UPGRADE_ROLLBACK.md index 2d405ad..f69bc00 100644 --- a/docs/UPGRADE_ROLLBACK.md +++ b/docs/UPGRADE_ROLLBACK.md @@ -1,7 +1,7 @@ # Upgrade ve rollback runbook'u -Bu runbook aynı PostgreSQL major sürümünde SQL Dashboard `1.1.0` ve repository -şema `0014` için güvenli yükseltme sınırını tanımlar. Üretimde image'ları sürüm +Bu runbook aynı PostgreSQL major sürümünde SQL Dashboard `1.1.1` ve repository +şema `0016` için güvenli yükseltme sınırını tanımlar. Üretimde image'ları sürüm ve digest ile sabitleyin. Önce staging restore provası yapın. > **Uyarı:** `dropdb`, `docker compose down --volumes` ve yanlış Compose proje @@ -53,18 +53,20 @@ docker compose exec -T repository-db \ -c "SELECT max(version),count(*) FROM advisor_migrations.schema_migrations" ``` -`1.1.0` için beklenen son sürüm `0014` olmalıdır. +`1.1.1` için beklenen son sürüm `0016` olmalıdır. ## 3. Servisleri aç, doğrula ve cutover yap ```bash -docker compose up -d collector join-snapshotter evaluator api web +docker compose up -d collector join-snapshotter evaluator query-metrics-snapshot-worker +# 1h ve 24h state satırları ready olduktan sonra: +docker compose up -d api web docker compose ps curl -fsS http://127.0.0.1:8000/api/v1/health bash scripts/verify.sh ``` -Health yanıtında uygulama `1.1.0`, migration `current=expected=0014` ve +Health yanıtında uygulama `1.1.1`, migration `current=expected=0016` ve `upToDate=true` bekleyin. Collector için yeni snapshot zamanı ilerlemeden ve web/API smoke testi geçmeden trafiği açmayın. Dış source bağlantılarını alias, database OID ve capability kapsamıyla yeniden doğrulayın. diff --git a/frontend/package-lock.json b/frontend/package-lock.json index 22a3c8a..a9a7e87 100644 --- a/frontend/package-lock.json +++ b/frontend/package-lock.json @@ -1,12 +1,12 @@ { "name": "postgresql-advisor-web", - "version": "1.1.0", + "version": "1.1.1", "lockfileVersion": 3, "requires": true, "packages": { "": { "name": "postgresql-advisor-web", - "version": "1.1.0", + "version": "1.1.1", "dependencies": { "@dagrejs/dagre": "3.0.0", "@vitejs/plugin-react": "6.0.4", diff --git a/frontend/package.json b/frontend/package.json index 2029c52..1319d36 100644 --- a/frontend/package.json +++ b/frontend/package.json @@ -1,7 +1,7 @@ { "name": "postgresql-advisor-web", "private": true, - "version": "1.1.0", + "version": "1.1.1", "type": "module", "scripts": { "dev": "vite --host 127.0.0.1", diff --git a/scripts/register-source.sh b/scripts/register-source.sh index 6fd3583..1ea3ab4 100755 --- a/scripts/register-source.sh +++ b/scripts/register-source.sh @@ -272,11 +272,11 @@ probe="$({ printf '%s\n' "$source_password"; } | docker compose exec -T reposito AND current_setting('\''pg_wait_sampling.profile_queries'\'') = '\''top'\'' AND current_setting('\''pg_wait_sampling.sample_cpu'\'') = '\''off'\'', COALESCE(has_function_privilege( - current_user, 'advisor_join.capture_and_reset()', 'EXECUTE' + current_user, '\''advisor_join.capture_and_reset()'\'', '\''EXECUTE'\'' ), false), - to_regprocedure('advisor_join.assert_outbox_within_limits()') IS NOT NULL, + to_regprocedure('\''advisor_join.assert_outbox_within_limits()'\'') IS NOT NULL, NOT COALESCE(has_function_privilege( - current_user, 'advisor_join.assert_outbox_within_limits()', 'EXECUTE' + current_user, '\''advisor_join.assert_outbox_within_limits()'\'', '\''EXECUTE'\'' ), false), NOT COALESCE(( SELECT has_function_privilege( diff --git a/scripts/verify.sh b/scripts/verify.sh index f829caa..89e0e16 100755 --- a/scripts/verify.sh +++ b/scripts/verify.sh @@ -978,6 +978,27 @@ collector_state="$(docker compose exec -T repository-db psql -U postgres -p 5433 [[ "$collector_state" == "HEALTHY|0" ]] || fail "Collector sagligi beklenmiyor: ${collector_state}" pass "Collector gecikmesi ve hata durumu saglikli" +docker compose exec -T query-metrics-snapshot-worker python - <<'PY' +import time + +from app.config import get_settings +from app.snapshot_worker import open_connection, refresh_snapshot + + +settings = get_settings() +with open_connection(settings) as connection: + for window in ("1h", "24h"): + for _ in range(30): + if refresh_snapshot(connection, window): + break + time.sleep(2) + else: + raise SystemExit( + f"{window} query metrics snapshot advisory lock could not be acquired" + ) +PY +pass "Dashboard 1h/24h persistent query snapshot'lari test deltasi sonrasi yenilendi" + curl -fsS -H 'X-Advisor-Role: analyst' \ "${api_url}/api/v1/queries?window=1h&pageSize=10" >"${verify_tmp_dir}/queries-authorized.json" curl -fsS "${api_url}/api/v1/queries?window=1h&pageSize=10" >"${verify_tmp_dir}/queries-viewer.json" diff --git a/sql/017_query_metrics_snapshots.sql b/sql/017_query_metrics_snapshots.sql new file mode 100644 index 0000000..c636c19 --- /dev/null +++ b/sql/017_query_metrics_snapshots.sql @@ -0,0 +1,80 @@ +\set ON_ERROR_STOP on + +-- Query metrics are expensive at ERP fingerprint cardinalities. Keep one +-- persistent, atomically replaceable materialized snapshot per supported +-- dashboard window. The refresh worker populates these outside the migration +-- transaction and uses REFRESH MATERIALIZED VIEW CONCURRENTLY after the first +-- population, so readers retain the previous complete generation. + +CREATE MATERIALIZED VIEW advisor.query_metrics_snapshot_1h AS +SELECT metrics.*, servers.alias AS server_alias +FROM advisor.query_metrics(interval '1 hour') AS metrics +LEFT JOIN "PoWA".powa_servers AS servers + ON servers.id = metrics.server_id +WITH NO DATA; + +CREATE MATERIALIZED VIEW advisor.query_metrics_snapshot_24h AS +SELECT metrics.*, servers.alias AS server_alias +FROM advisor.query_metrics(interval '24 hours') AS metrics +LEFT JOIN "PoWA".powa_servers AS servers + ON servers.id = metrics.server_id +WITH NO DATA; + +CREATE MATERIALIZED VIEW advisor.query_metrics_snapshot_7d AS +SELECT metrics.*, servers.alias AS server_alias +FROM advisor.query_metrics(interval '7 days') AS metrics +LEFT JOIN "PoWA".powa_servers AS servers + ON servers.id = metrics.server_id +WITH NO DATA; + +CREATE MATERIALIZED VIEW advisor.query_metrics_snapshot_30d AS +SELECT metrics.*, servers.alias AS server_alias +FROM advisor.query_metrics(interval '30 days') AS metrics +LEFT JOIN "PoWA".powa_servers AS servers + ON servers.id = metrics.server_id +WITH NO DATA; + +CREATE UNIQUE INDEX query_metrics_snapshot_1h_identity_idx + ON advisor.query_metrics_snapshot_1h + (server_id, database_id, query_id, user_id) NULLS NOT DISTINCT; +CREATE UNIQUE INDEX query_metrics_snapshot_24h_identity_idx + ON advisor.query_metrics_snapshot_24h + (server_id, database_id, query_id, user_id) NULLS NOT DISTINCT; +CREATE UNIQUE INDEX query_metrics_snapshot_7d_identity_idx + ON advisor.query_metrics_snapshot_7d + (server_id, database_id, query_id, user_id) NULLS NOT DISTINCT; +CREATE UNIQUE INDEX query_metrics_snapshot_30d_identity_idx + ON advisor.query_metrics_snapshot_30d + (server_id, database_id, query_id, user_id) NULLS NOT DISTINCT; + +CREATE TABLE advisor.query_metrics_snapshot_state ( + window_key text PRIMARY KEY + CHECK (window_key IN ('1h', '24h', '7d', '30d')), + status text NOT NULL DEFAULT 'pending' + CHECK (status IN ('pending', 'refreshing', 'ready', 'failed')), + refresh_started_at timestamptz, + refreshed_at timestamptz, + refresh_duration_ms bigint + CHECK (refresh_duration_ms IS NULL OR refresh_duration_ms >= 0), + row_count bigint + CHECK (row_count IS NULL OR row_count >= 0), + last_error text, + updated_at timestamptz NOT NULL DEFAULT clock_timestamp() +); + +INSERT INTO advisor.query_metrics_snapshot_state (window_key) +VALUES ('1h'), ('24h'), ('7d'), ('30d'); + +REVOKE ALL ON advisor.query_metrics_snapshot_state FROM PUBLIC; +GRANT SELECT, INSERT, UPDATE ON advisor.query_metrics_snapshot_state + TO advisor_api; + +GRANT SELECT, MAINTAIN ON + advisor.query_metrics_snapshot_1h, + advisor.query_metrics_snapshot_24h, + advisor.query_metrics_snapshot_7d, + advisor.query_metrics_snapshot_30d +TO advisor_api; + +COMMENT ON TABLE advisor.query_metrics_snapshot_state IS + 'Persistent refresh status for precomputed dashboard query-metrics windows.'; diff --git a/sql/018_global_trend_snapshots.sql b/sql/018_global_trend_snapshots.sql new file mode 100644 index 0000000..a0410cd --- /dev/null +++ b/sql/018_global_trend_snapshots.sql @@ -0,0 +1,211 @@ +\set ON_ERROR_STOP on + +-- Overview trend aggregation is also expensive at ERP cardinalities. Compute +-- database-level trends once, then derive server and fleet totals from those +-- small bucket sets. This avoids rescanning the same query histories for every +-- dashboard scope. + +CREATE MATERIALIZED VIEW advisor.global_trend_snapshot_1h AS +WITH database_trends AS MATERIALIZED ( + SELECT + database.srvid AS server_id, + database.oid AS database_id, + trend.bucket_at AS timestamp, + trend.total_exec_time_ms, + trend.calls + FROM "PoWA".powa_databases AS database + CROSS JOIN LATERAL advisor.query_trend( + now() - interval '1 hour', + interval '5 minutes', + database.srvid, + database.oid + ) AS trend + WHERE database.datname <> 'powa' +), all_scopes AS ( + SELECT * FROM database_trends + UNION ALL + SELECT + server_id, + NULL::oid, + timestamp, + sum(total_exec_time_ms)::double precision, + sum(calls)::bigint + FROM database_trends + GROUP BY server_id, timestamp + UNION ALL + SELECT + NULL::integer, + NULL::oid, + timestamp, + sum(total_exec_time_ms)::double precision, + sum(calls)::bigint + FROM database_trends + GROUP BY timestamp +) +SELECT * FROM all_scopes +WITH NO DATA; + +CREATE MATERIALIZED VIEW advisor.global_trend_snapshot_24h AS +WITH database_trends AS MATERIALIZED ( + SELECT + database.srvid AS server_id, + database.oid AS database_id, + trend.bucket_at AS timestamp, + trend.total_exec_time_ms, + trend.calls + FROM "PoWA".powa_databases AS database + CROSS JOIN LATERAL advisor.query_trend( + now() - interval '24 hours', + interval '1 hour', + database.srvid, + database.oid + ) AS trend + WHERE database.datname <> 'powa' +), all_scopes AS ( + SELECT * FROM database_trends + UNION ALL + SELECT + server_id, + NULL::oid, + timestamp, + sum(total_exec_time_ms)::double precision, + sum(calls)::bigint + FROM database_trends + GROUP BY server_id, timestamp + UNION ALL + SELECT + NULL::integer, + NULL::oid, + timestamp, + sum(total_exec_time_ms)::double precision, + sum(calls)::bigint + FROM database_trends + GROUP BY timestamp +) +SELECT * FROM all_scopes +WITH NO DATA; + +CREATE MATERIALIZED VIEW advisor.global_trend_snapshot_7d AS +WITH database_trends AS MATERIALIZED ( + SELECT + database.srvid AS server_id, + database.oid AS database_id, + trend.bucket_at AS timestamp, + trend.total_exec_time_ms, + trend.calls + FROM "PoWA".powa_databases AS database + CROSS JOIN LATERAL advisor.query_trend( + now() - interval '7 days', + interval '6 hours', + database.srvid, + database.oid + ) AS trend + WHERE database.datname <> 'powa' +), all_scopes AS ( + SELECT * FROM database_trends + UNION ALL + SELECT + server_id, + NULL::oid, + timestamp, + sum(total_exec_time_ms)::double precision, + sum(calls)::bigint + FROM database_trends + GROUP BY server_id, timestamp + UNION ALL + SELECT + NULL::integer, + NULL::oid, + timestamp, + sum(total_exec_time_ms)::double precision, + sum(calls)::bigint + FROM database_trends + GROUP BY timestamp +) +SELECT * FROM all_scopes +WITH NO DATA; + +CREATE MATERIALIZED VIEW advisor.global_trend_snapshot_30d AS +WITH database_trends AS MATERIALIZED ( + SELECT + database.srvid AS server_id, + database.oid AS database_id, + trend.bucket_at AS timestamp, + trend.total_exec_time_ms, + trend.calls + FROM "PoWA".powa_databases AS database + CROSS JOIN LATERAL advisor.query_trend( + now() - interval '30 days', + interval '1 day', + database.srvid, + database.oid + ) AS trend + WHERE database.datname <> 'powa' +), all_scopes AS ( + SELECT * FROM database_trends + UNION ALL + SELECT + server_id, + NULL::oid, + timestamp, + sum(total_exec_time_ms)::double precision, + sum(calls)::bigint + FROM database_trends + GROUP BY server_id, timestamp + UNION ALL + SELECT + NULL::integer, + NULL::oid, + timestamp, + sum(total_exec_time_ms)::double precision, + sum(calls)::bigint + FROM database_trends + GROUP BY timestamp +) +SELECT * FROM all_scopes +WITH NO DATA; + +CREATE UNIQUE INDEX global_trend_snapshot_1h_identity_idx + ON advisor.global_trend_snapshot_1h + (server_id, database_id, timestamp) NULLS NOT DISTINCT; +CREATE UNIQUE INDEX global_trend_snapshot_24h_identity_idx + ON advisor.global_trend_snapshot_24h + (server_id, database_id, timestamp) NULLS NOT DISTINCT; +CREATE UNIQUE INDEX global_trend_snapshot_7d_identity_idx + ON advisor.global_trend_snapshot_7d + (server_id, database_id, timestamp) NULLS NOT DISTINCT; +CREATE UNIQUE INDEX global_trend_snapshot_30d_identity_idx + ON advisor.global_trend_snapshot_30d + (server_id, database_id, timestamp) NULLS NOT DISTINCT; + +CREATE TABLE advisor.global_trend_snapshot_state ( + window_key text PRIMARY KEY + CHECK (window_key IN ('1h', '24h', '7d', '30d')), + status text NOT NULL DEFAULT 'pending' + CHECK (status IN ('pending', 'refreshing', 'ready', 'failed')), + refresh_started_at timestamptz, + refreshed_at timestamptz, + refresh_duration_ms bigint + CHECK (refresh_duration_ms IS NULL OR refresh_duration_ms >= 0), + row_count bigint + CHECK (row_count IS NULL OR row_count >= 0), + last_error text, + updated_at timestamptz NOT NULL DEFAULT clock_timestamp() +); + +INSERT INTO advisor.global_trend_snapshot_state (window_key) +VALUES ('1h'), ('24h'), ('7d'), ('30d'); + +REVOKE ALL ON advisor.global_trend_snapshot_state FROM PUBLIC; +GRANT SELECT, INSERT, UPDATE ON advisor.global_trend_snapshot_state + TO advisor_api; + +GRANT SELECT, MAINTAIN ON + advisor.global_trend_snapshot_1h, + advisor.global_trend_snapshot_24h, + advisor.global_trend_snapshot_7d, + advisor.global_trend_snapshot_30d +TO advisor_api; + +COMMENT ON TABLE advisor.global_trend_snapshot_state IS + 'Persistent refresh status for precomputed overview trends and scopes.'; diff --git a/sql/repository-migrations.manifest b/sql/repository-migrations.manifest index 5876ac4..7438b52 100644 --- a/sql/repository-migrations.manifest +++ b/sql/repository-migrations.manifest @@ -13,3 +13,5 @@ 0012|query_metrics_scale_guard|014_query_metrics_scale_guard.sql|a4c1484f00a8d2881e25f8f7fa6b2616bee475ed375a32782b9940b6ae2495ec 0013|query_trend_performance|015_query_trend_performance.sql|c600e3bffb01c4f7efe19b693039e8b17b338526e00ae5f429cf23fbcc377309 0014|product_scope_optimize_release|016_product_scope_optimize_release.sql|cf33caad1d4c5e39eb1144b47eaed5cf53448763636f79a069ef40035dcf76c1 +0015|query_metrics_snapshots|017_query_metrics_snapshots.sql|c6473663fdb41e9e77c25ce365fe4c3fa8f763cb65b9292b063b4aa910e87ea5 +0016|global_trend_snapshots|018_global_trend_snapshots.sql|99af76159349658100894ab320ad7b6c59e5ccb53290642a8a819f8a80d73320