This commit is contained in:
JiriUhlir
2026-06-18 11:58:23 +02:00
parent 284e013753
commit 6934f22253
18 changed files with 1367 additions and 17 deletions
View File
+180
View File
@@ -0,0 +1,180 @@
"""Google Analytics client.
Thin proxy over the GA4 Data API and Admin API. Request/response bodies are
forwarded as-is so callers keep the full flexibility of Google's API; this
module only handles authentication (Bearer token), the base URL, and error
mapping.
Authentication (per chosen model, token wins over service account):
* If ``X-GA-Access-Token`` was supplied, it is used directly.
* Otherwise a short-lived access token is minted from the service-account
JSON via google-auth and cached in-memory (keyed by key id + scope) until
shortly before it expires. The key material is never written to disk or log.
"""
from __future__ import annotations
import hashlib
import threading
import time
from typing import Any
import httpx
from fastapi.concurrency import run_in_threadpool
from .. import config
from ..credentials import GaCredentials
from ..errors import MissingCredentialsError, UpstreamError
from ..logging_config import get_logger
logger = get_logger(__name__)
# In-memory access-token cache: { cache_key: (token, expiry_epoch_seconds) }.
# Memory only - mirrors the stateless design (no secret ever persisted).
_token_cache: dict[str, tuple[str, float]] = {}
_token_lock = threading.Lock()
# Refresh a minted token this many seconds before its real expiry.
_EXPIRY_SKEW = 60.0
def _mint_token_sync(info: dict, scope: str) -> tuple[str, float]:
"""Mint an OAuth2 access token from a service-account key (blocking)."""
# Imported lazily so the module imports even if google-auth is missing,
# and so the dependency is only needed when service-account auth is used.
from google.auth.transport.requests import Request
from google.oauth2 import service_account
try:
creds = service_account.Credentials.from_service_account_info(
info, scopes=[scope]
)
except (ValueError, KeyError) as exc:
raise MissingCredentialsError(
f"X-GA-Credentials is not a usable service-account key: {exc}"
) from exc
try:
creds.refresh(Request())
except Exception as exc: # google.auth.exceptions.RefreshError and friends
# Surface as upstream auth failure - do NOT log the key material.
raise UpstreamError(
f"Failed to obtain Google access token from service account: {exc}",
status=401,
) from exc
expiry = creds.expiry.timestamp() if creds.expiry else (time.time() + 3600)
return creds.token, expiry
def _cache_key(info: dict, scope: str) -> str:
# Identify a key by its private_key_id + client_email + scope. Hashed so the
# raw identifiers never sit in a dict key we might later log.
raw = f"{info.get('private_key_id', '')}|{info.get('client_email', '')}|{scope}"
return hashlib.sha256(raw.encode("utf-8")).hexdigest()
async def _bearer_token(creds: GaCredentials) -> str:
if creds.access_token:
return creds.access_token
if creds.service_account_info is None:
# get_ga_credentials guarantees one of the two, but be defensive.
raise MissingCredentialsError(
"No GA access token and no service-account credentials available."
)
scope = config.GA_SCOPE
key = _cache_key(creds.service_account_info, scope)
now = time.time()
with _token_lock:
cached = _token_cache.get(key)
if cached and cached[1] - _EXPIRY_SKEW > now:
return cached[0]
# Mint outside the lock (network call); google-auth is blocking, so offload
# it to a thread to avoid stalling the event loop.
token, expiry = await run_in_threadpool(
_mint_token_sync, creds.service_account_info, scope
)
with _token_lock:
_token_cache[key] = (token, expiry)
return token
class GoogleAnalyticsClient:
"""Authenticated HTTP client for the GA4 Data and Admin APIs."""
def __init__(self, creds: GaCredentials) -> None:
self._creds = creds
async def _request(
self,
method: str,
base_url: str,
path: str,
*,
params: dict | None = None,
json_body: Any | None = None,
) -> Any:
token = await _bearer_token(self._creds)
headers = {"Authorization": f"Bearer {token}"}
if self._creds.quota_project:
headers["x-goog-user-project"] = self._creds.quota_project
url = f"{base_url}{path}"
try:
async with httpx.AsyncClient(
timeout=config.HTTP_TIMEOUT_SECONDS
) as client:
resp = await client.request(
method, url, params=params, json=json_body, headers=headers
)
except httpx.TimeoutException as exc:
raise UpstreamError(
"Google Analytics request timed out.", status=504
) from exc
except httpx.HTTPError as exc:
raise UpstreamError(
f"Google Analytics is unreachable: {exc}", status=502
) from exc
return _parse_google_response(resp)
# --- Data API -------------------------------------------------------------
async def data_post(self, path: str, body: Any) -> Any:
return await self._request(
"POST", config.GA_DATA_BASE_URL, path, json_body=body
)
async def data_get(self, path: str, params: dict | None = None) -> Any:
return await self._request(
"GET", config.GA_DATA_BASE_URL, path, params=params
)
# --- Admin API ------------------------------------------------------------
async def admin_get(self, path: str, params: dict | None = None) -> Any:
return await self._request(
"GET", config.GA_ADMIN_BASE_URL, path, params=params
)
def _parse_google_response(resp: httpx.Response) -> Any:
try:
payload = resp.json()
except ValueError:
payload = {"raw": resp.text}
if resp.is_success:
return payload
# Google returns {"error": {"code", "message", "status", ...}}.
message = "Google Analytics API error"
if isinstance(payload, dict) and isinstance(payload.get("error"), dict):
message = payload["error"].get("message", message)
raise UpstreamError(
message,
status=502 if resp.status_code >= 500 else resp.status_code,
upstream_status=resp.status_code,
body=payload,
)
+193
View File
@@ -0,0 +1,193 @@
"""Sklik (Seznam) "Drak" JSON API client.
Protocol (verified against the official seznam/api-examples JSON example):
* Endpoint: ``{SKLIK_BASE_URL}/{method}`` e.g. .../drak/json/v5/campaigns.list
* HTTP POST, body = a JSON ARRAY of positional arguments.
* ``client.loginByToken`` takes the API token as its single argument and
returns ``{"status":200,"session":"...",...}``.
* Every authenticated method takes the user struct ``{"session": ...}``
(optionally ``"userId"``) as its FIRST argument, followed by the method's
own arguments.
* Every response is an object containing ``status`` (HTTP-style int),
``statusMessage``, a refreshed ``session``, plus method-specific data.
The proxy is stateless: it logs in with ``X-Sklik-Token`` per request to obtain
a session, then performs the requested call. The token and session are never
logged.
"""
from __future__ import annotations
from typing import Any
import httpx
from .. import config
from ..errors import UpstreamError
from ..logging_config import get_logger
logger = get_logger(__name__)
# Sklik report data is paginated; readReport is called with an offset/limit
# window until all rows are fetched. Keep the page size conservative.
_REPORT_PAGE_LIMIT = 100
# Hard stop so a misbehaving upstream can't loop forever.
_REPORT_MAX_PAGES = 1000
class SklikClient:
"""Performs JSON-RPC calls against the Sklik Drak API."""
def __init__(self, token: str, user_id: int | None = None) -> None:
self._token = token
self._user_id = user_id
self._session: str | None = None
async def __aenter__(self) -> "SklikClient":
self._http = httpx.AsyncClient(timeout=config.HTTP_TIMEOUT_SECONDS)
return self
async def __aexit__(self, *exc: Any) -> None:
await self._http.aclose()
async def _call(self, method: str, args: list[Any]) -> dict:
"""Low-level: POST a JSON array of args to ``/{method}``."""
url = f"{config.SKLIK_BASE_URL}/{method}"
try:
resp = await self._http.post(url, json=args)
except httpx.TimeoutException as exc:
raise UpstreamError(
f"Sklik request timed out ({method}).", status=504
) from exc
except httpx.HTTPError as exc:
raise UpstreamError(
f"Sklik is unreachable ({method}): {exc}", status=502
) from exc
try:
payload = resp.json()
except ValueError as exc:
raise UpstreamError(
f"Sklik returned a non-JSON response ({method}).",
status=502,
upstream_status=resp.status_code,
body={"raw": resp.text},
) from exc
if not isinstance(payload, dict):
raise UpstreamError(
f"Unexpected Sklik response shape ({method}).",
status=502,
body=payload,
)
status = payload.get("status")
# Sklik conveys business errors in the body with an HTTP-style status.
# 200 OK, 206 partially OK, 301 "user is serviced" are all acceptable.
if status not in (200, 206, 301):
raise UpstreamError(
payload.get("statusMessage", f"Sklik error on {method}."),
status=400 if isinstance(status, int) and 400 <= status < 500 else 502,
upstream_status=status if isinstance(status, int) else None,
body=payload,
)
# Refresh our session from every response (Sklik rotates it).
new_session = payload.get("session")
if isinstance(new_session, str) and new_session:
self._session = new_session
return payload
async def login(self) -> dict:
"""Exchange the API token for a session. Idempotent per client."""
payload = await self._call("client.loginByToken", [self._token])
if not self._session:
raise UpstreamError(
"Sklik login succeeded but returned no session.",
status=502,
body=payload,
)
return payload
def _user_struct(self) -> dict:
user: dict[str, Any] = {"session": self._session}
if self._user_id is not None:
user["userId"] = self._user_id
return user
async def call(self, method: str, args: list[Any] | None = None) -> dict:
"""Authenticated call: prepends the user/session struct to ``args``.
Logs in first if no session is held yet. ``method`` is e.g.
``campaigns.list``; ``args`` are the method arguments AFTER the user
struct.
"""
if method == "client.loginByToken":
# Login is handled by login(); never forward the bare token here.
raise UpstreamError(
"client.loginByToken cannot be called directly; the proxy "
"manages the session.",
status=400,
)
if not self._session:
await self.login()
full_args = [self._user_struct()] + list(args or [])
return await self._call(method, full_args)
async def fetch_report(
self, entity: str, report_args: list[Any]
) -> dict:
"""Create a stats report for ``entity`` then read all of its rows.
``entity`` is e.g. ``campaigns``/``groups``/``ads``/``keywords``.
Calls ``{entity}.createReport`` with ``report_args`` (the restriction +
display-options structs), then pages through ``{entity}.readReport``
until every row is collected.
"""
created = await self.call(f"{entity}.createReport", report_args)
report_id = created.get("reportId")
if not report_id:
raise UpstreamError(
f"{entity}.createReport returned no reportId.",
status=502,
body=created,
)
total = created.get("totalCount", 0)
rows: list[Any] = []
offset = 0
pages = 0
while True:
page = await self.call(
f"{entity}.readReport",
[
report_id,
{
"offset": offset,
"limit": _REPORT_PAGE_LIMIT,
"allowEmptyStatistics": False,
},
],
)
batch = page.get("report") or []
rows.extend(batch)
pages += 1
offset += _REPORT_PAGE_LIMIT
if len(batch) < _REPORT_PAGE_LIMIT:
break
if pages >= _REPORT_MAX_PAGES:
logger.warning(
"Sklik %s.readReport hit the %d-page safety cap (collected "
"%d rows); result may be truncated.",
entity,
_REPORT_MAX_PAGES,
len(rows),
)
break
return {
"reportId": report_id,
"totalCount": total,
"returnedCount": len(rows),
"truncated": pages >= _REPORT_MAX_PAGES,
"report": rows,
}