This commit is contained in:
2026-07-09 02:52:54 +02:00
parent 818568c609
commit 48bb6c2ec2
95 changed files with 6115 additions and 267 deletions
+1
View File
@@ -4,6 +4,7 @@ CATEGORY_BY_PREFIX: dict[str, str] = {
"auth": "auth",
"profile": "account",
"follow": "social",
"award": "social",
"relation": "social",
"push": "push",
"notification": "notification",
+10
View File
@@ -64,6 +64,16 @@ class BackgroundQueue:
if inline:
self._inline += 1
logger.warning("background task %s failed: %s", getattr(fn, "__name__", fn), exc)
finally:
self._release_db_lock()
def _release_db_lock(self) -> None:
try:
from devplacepy.database import refresh_snapshot
refresh_snapshot()
except Exception as exc:
logger.warning("background queue could not release db lock: %s", exc)
async def start(self) -> None:
if self.running:
@@ -31,6 +31,25 @@ ADMIN_ACTIONS: tuple[Action, ...] = (
params=(query("top_n", "How many top authors to include (1-50)."),),
requires_admin=True,
),
Action(
name="admin_statistics",
method="GET",
path="/admin/statistics/data",
summary="Platform statistics tab data with trends (admin only)",
description=(
"Returns JSON for one statistics tab: KPI cards with period-over-period deltas, "
"time-series points for line charts, breakdown tables, and highlight metrics. "
"Tabs: overview, visitors, members, content, engagement, social, ai, devii, "
"services, containers, game, awards, moderation, tools, storage."
),
params=(
query("tab", "Tab key (default overview)."),
query("hours", "Lookback window in hours (24, 168, 720, 2160, or 0 for all time)."),
query("compare", "Include previous-period comparison (1 or 0, default 1)."),
query("top_n", "Rows in breakdown tables (default 10)."),
),
requires_admin=True,
),
Action(
name="ai_usage",
method="GET",
@@ -305,6 +324,18 @@ ADMIN_ACTIONS: tuple[Action, ...] = (
params=(path("uid", "Attachment uid to purge."), confirm()),
requires_admin=True,
),
Action(
name="admin_revoke_award",
method="POST",
path="/admin/awards/{uid}/revoke",
summary="Revoke a published award (admin only)",
description=(
"Soft-deletes the award and linked attachments, then recomputes receiver stats. "
"Confirmation is required in the UI; Devii should confirm before calling."
),
params=(path("uid", "Award uid to revoke."), confirm()),
requires_admin=True,
),
Action(
name="admin_reset_guest_ai_quota",
method="POST",
@@ -27,7 +27,7 @@ GATEWAY_ACTIONS: tuple[Action, ...] = (
summary="Create or update a gateway provider (admin only)",
description=(
"Adds or updates a named upstream provider. base_url is the OpenAI-compatible "
"chat-completions endpoint; the embeddings endpoint is derived from it."
"chat-completions endpoint; the embeddings and images endpoints are derived from it."
),
handler="http",
requires_admin=True,
@@ -60,7 +60,8 @@ GATEWAY_ACTIONS: tuple[Action, ...] = (
summary="List OpenAI gateway model routes (admin only)",
description=(
"Returns JSON: every source-model route with its provider, target model, kind "
"(chat/embed), optional vision model, context window, and per-model pricing economy."
"(chat/embed/image), optional vision model, context window, and per-model pricing economy "
"(including any tiered/off-peak pricing configured on it)."
),
handler="http",
requires_admin=True,
@@ -73,9 +74,14 @@ GATEWAY_ACTIONS: tuple[Action, ...] = (
summary="Create or update a gateway model route (admin only)",
description=(
"Maps a requested source_model onto a provider + target_model, each with its own "
"pricing. kind is 'chat' or 'embed'. A vision_model adds image-to-text augmentation "
"for chat routes. Prices are USD per 1,000,000 tokens; chat uses cache-hit/cache-miss/"
"output, embeddings use input, vision uses input/output."
"pricing. kind is 'chat', 'embed', or 'image'. A vision_model adds image-to-text augmentation "
"for chat routes. Prices are USD per 1,000,000 tokens for chat/embed/vision; image routes "
"use price_input_per_m as a flat USD per generated image. Optionally, a route can also "
"charge a different (tier-2) rate once the request's input tokens exceed "
"context_tier_threshold_tokens, and/or apply a percentage discount during a fixed "
"UTC off-peak window - leave the tier2/off-peak fields unset to keep the flat rates "
"above at all times. off_peak_start_minute and off_peak_end_minute must be set "
"together (both or neither) or the call is rejected."
),
handler="http",
requires_admin=True,
@@ -83,14 +89,22 @@ GATEWAY_ACTIONS: tuple[Action, ...] = (
body("source_model", "Model name clients request.", required=True),
body("provider", "Provider name, or blank for the default upstream."),
body("target_model", "Model name sent upstream.", required=True),
body("kind", "Route kind: 'chat' or 'embed'."),
body("kind", "Route kind: 'chat', 'embed', or 'image'."),
body("vision_provider", "Provider for image description, or blank for the route provider."),
body("vision_model", "Vision model name (blank disables the merge)."),
Param(name="context_window", location="body", description="Max context tokens (0 = unknown).", required=False, type="integer"),
Param(name="price_cache_hit_per_m", location="body", description="USD per 1M cache-hit input tokens.", required=False, type="number"),
Param(name="price_cache_miss_per_m", location="body", description="USD per 1M cache-miss input tokens.", required=False, type="number"),
Param(name="price_output_per_m", location="body", description="USD per 1M output tokens.", required=False, type="number"),
Param(name="price_input_per_m", location="body", description="USD per 1M input tokens (embed/vision).", required=False, type="number"),
Param(name="price_input_per_m", location="body", description="USD per 1M input tokens (embed/vision), or flat USD per image for image routes.", required=False, type="number"),
Param(name="context_tier_threshold_tokens", location="body", description="Input tokens above which tier-2 rates apply (0 disables tiering).", required=False, type="integer"),
Param(name="price_cache_hit_per_m_tier2", location="body", description="Tier-2 USD per 1M cache-hit input tokens (unset = keep tier-1 rate above threshold).", required=False, type="number"),
Param(name="price_cache_miss_per_m_tier2", location="body", description="Tier-2 USD per 1M cache-miss input tokens (unset = keep tier-1 rate above threshold).", required=False, type="number"),
Param(name="price_output_per_m_tier2", location="body", description="Tier-2 USD per 1M output tokens (unset = keep tier-1 rate above threshold).", required=False, type="number"),
Param(name="price_input_per_m_tier2", location="body", description="Tier-2 USD per 1M input tokens, embed/vision (unset = keep tier-1 rate above threshold).", required=False, type="number"),
Param(name="off_peak_start_minute", location="body", description="Off-peak window start, UTC minutes since midnight (0-1439). Must be set together with off_peak_end_minute.", required=False, type="integer"),
Param(name="off_peak_end_minute", location="body", description="Off-peak window end, UTC minutes since midnight (0-1439). A value less than the start wraps past midnight.", required=False, type="integer"),
Param(name="off_peak_discount_pct", location="body", description="Percentage discount (0-100) applied to the active tier's rates during the off-peak window.", required=False, type="number"),
Param(name="is_active", location="body", description="Whether the route is active ('1' or '0').", required=False, type="boolean"),
),
),
@@ -61,6 +61,21 @@ PROFILE_ACTIONS: tuple[Action, ...] = (
),
),
),
Action(
name="give_award",
method="POST",
path="/profile/{username}/award",
summary="Give another member an award on their profile",
description=(
"Creates a pending award row and enqueues image generation billed to the giver's "
"API key. The receiver is notified when generation completes. Cannot award yourself; "
"blocked pairs are rejected; cooldowns apply."
),
params=(
path("username", "Receiver username."),
body("description", "Short award message (1-125 characters).", required=True),
),
),
Action(
name="regenerate_avatar",
method="POST",
@@ -18,7 +18,7 @@ def arg(
TYPES = (
"One of: comment, reply, mention, vote, follow, message, badge, level, issue."
"One of: comment, reply, mention, vote, follow, message, badge, level, issue, award."
)
NOTIFICATION_ACTIONS: tuple[Action, ...] = (
+37 -60
View File
@@ -21,32 +21,6 @@ from ..text import html_to_text
logger = logging.getLogger("devii.fetch")
CHROME_VERSION = "131"
CHROME_VERSION_FULL = "131.0.0.0"
STEALTH_HEADERS = {
"User-Agent": (
"Mozilla/5.0 (X11; Linux x86_64) AppleWebKit/537.36 "
f"(KHTML, like Gecko) Chrome/{CHROME_VERSION_FULL} Safari/537.36"
),
"Accept": (
"text/html,application/xhtml+xml,application/xml;q=0.9,image/avif,image/webp,"
"image/apng,*/*;q=0.8,application/signed-exchange;v=b3;q=0.7"
),
"Accept-Language": "en-US,en;q=0.9",
"sec-ch-ua": (
f'"Google Chrome";v="{CHROME_VERSION}", "Chromium";v="{CHROME_VERSION}", '
'"Not_A Brand";v="24"'
),
"sec-ch-ua-mobile": "?0",
"sec-ch-ua-platform": '"Linux"',
"Sec-Fetch-Dest": "document",
"Sec-Fetch-Mode": "navigate",
"Sec-Fetch-Site": "none",
"Sec-Fetch-User": "?1",
"Upgrade-Insecure-Requests": "1",
"Cache-Control": "max-age=0",
"DNT": "1",
}
TITLE = re.compile(r"<title[^>]*>(.*?)</title>", re.IGNORECASE | re.DOTALL)
MIN_FETCH_CHARS = 1000
ALLOWED_METHODS = ("GET", "POST", "PUT", "PATCH", "DELETE", "HEAD", "OPTIONS")
@@ -230,9 +204,7 @@ class FetchController:
raise_5xx: bool = False,
) -> tuple[str, str, str, int, dict[str, str]]:
limit = self._settings.fetch_max_bytes
merged_headers = dict(STEALTH_HEADERS)
if headers:
merged_headers.update({str(k): str(v) for k, v in headers.items()})
merged_headers = {str(k): str(v) for k, v in headers.items()} if headers else {}
request_kwargs: dict[str, Any] = {}
if json_body is not None:
request_kwargs["json"] = json_body
@@ -240,43 +212,48 @@ class FetchController:
request_kwargs["data"] = form
elif content is not None:
request_kwargs["content"] = content
async def _once(client: httpx.AsyncClient, target: str) -> tuple[str, str, str, int, dict[str, str]]:
async with client.stream(method, target, **request_kwargs) as response:
if raise_5xx and response.status_code >= 500:
raise UpstreamError(
f"Server returned {response.status_code}.",
status=response.status_code,
url=target,
)
chunks: list[bytes] = []
total = 0
async for chunk in response.aiter_bytes():
chunks.append(chunk)
total += len(chunk)
if total >= limit:
break
raw = b"".join(chunks)[:limit]
encoding = response.encoding or "utf-8"
try:
body = raw.decode(encoding, errors="replace")
except LookupError:
body = raw.decode("utf-8", errors="replace")
content_type = response.headers.get("content-type", "").lower()
response_headers = {key: value for key, value in response.headers.items()}
return body, str(response.url), content_type, response.status_code, response_headers
try:
async with stealth.stealth_async_client(
headers=merged_headers,
follow_redirects=True,
timeout=self._settings.fetch_timeout_seconds,
) as client:
async with client.stream(method, url, **request_kwargs) as response:
if raise_5xx and response.status_code >= 500:
raise UpstreamError(
f"Server returned {response.status_code}.",
status=response.status_code,
url=url,
)
chunks: list[bytes] = []
total = 0
async for chunk in response.aiter_bytes():
chunks.append(chunk)
total += len(chunk)
if total >= limit:
break
raw = b"".join(chunks)[:limit]
encoding = response.encoding or "utf-8"
try:
body = raw.decode(encoding, errors="replace")
except LookupError:
body = raw.decode("utf-8", errors="replace")
content_type = response.headers.get("content-type", "").lower()
response_headers = {
key: value for key, value in response.headers.items()
}
return (
body,
str(response.url),
content_type,
response.status_code,
response_headers,
)
result = await _once(client, url)
if method == "GET" and "html" in result[2]:
gate_url = stealth.detect_consent_gate(result[0])
if gate_url:
try:
await client.get(gate_url)
result = await _once(client, url)
except httpx.HTTPError as exc:
logger.debug("consent gate follow-up failed for %s: %s", url, exc)
return result
except httpx.TimeoutException as exc:
raise NetworkError(f"Request timed out fetching {url}", url=url) from exc
except httpx.HTTPError as exc:
+207
View File
@@ -0,0 +1,207 @@
# retoor <retoor@molodetz.nl>
import asyncio
import base64
import json
import logging
from datetime import datetime, timezone
from io import BytesIO
from devplacepy import stealth
from devplacepy.attachments import link_attachments, store_attachment
from devplacepy.awards.images import enforce_rgba_png, resize_award_png
from devplacepy.config import (
AWARD_GENERATION_TIMEOUT_SECONDS,
AWARD_IMAGE_MODEL_DEFAULT,
AWARD_IMAGE_PROMPT_DEFAULT,
AWARD_IMAGE_SIZE_DEFAULT,
INTERNAL_BASE_URL,
)
from devplacepy.database import add_award_usage, get_award_usage, get_setting, get_table
from devplacepy.database.awards import recompute_user_award_stats
from devplacepy.services.base import ConfigField
from devplacepy.services.correction import new_usage_totals
from devplacepy.services.jobs.base import JobService
from devplacepy.services.openai_gateway.usage import accumulate_usage, usage_metric_cards
from devplacepy.utils import create_notification
logger = logging.getLogger(__name__)
class AwardService(JobService):
kind = "award"
title = "Awards"
description = (
"Generates award images off the request path, stores PNG attachments, and "
"notifies receivers when ready."
)
def __init__(self):
super().__init__(name="award", interval_seconds=2)
self.config_fields = list(self.config_fields) + [
ConfigField(
"award_give_cooldown_hours",
"Give cooldown (hours)",
type="int",
default=24,
minimum=1,
maximum=168,
group="Awards",
),
ConfigField(
"award_receive_cooldown_hours",
"Receive cooldown (hours)",
type="int",
default=24,
minimum=1,
maximum=168,
group="Awards",
),
ConfigField(
"award_display_hours",
"Prominent display (hours)",
type="int",
default=24,
minimum=1,
maximum=168,
group="Awards",
),
ConfigField(
"award_image_model",
"Image model",
type="str",
default=AWARD_IMAGE_MODEL_DEFAULT,
group="Awards",
),
ConfigField(
"award_image_prompt",
"Image prompt",
type="text",
default=AWARD_IMAGE_PROMPT_DEFAULT,
group="Awards",
),
]
async def process(self, job: dict) -> dict:
totals = new_usage_totals()
try:
result = await asyncio.to_thread(self._generate, job, totals)
except Exception as exc:
self._audit_failed(job, str(exc))
raise
if totals["calls"]:
add_award_usage(totals)
return result
def _generate(self, job: dict, totals: dict) -> dict:
from devplacepy.services.audit import record as audit
payload = job.get("payload") or {}
award_uid = str(payload.get("award_uid", ""))
giver_uid = str(payload.get("giver_uid", ""))
receiver_uid = str(payload.get("receiver_uid", ""))
description = str(payload.get("description", ""))
api_key = str(payload.get("api_key", ""))
awards = get_table("awards")
row = awards.find_one(uid=award_uid)
if not row or row.get("deleted_at") or row.get("generated_at"):
return {"award_uid": award_uid, "skipped": True}
source_png = self._fetch_image(api_key, description, totals)
source_png = enforce_rgba_png(source_png)
uid512 = store_attachment(
resize_award_png(source_png, 512), "award-512.png", receiver_uid
)["uid"]
uid256 = store_attachment(
resize_award_png(source_png, 256), "award-256.png", giver_uid
)["uid"]
uid64 = store_attachment(
resize_award_png(source_png, 64), "award-64.png", giver_uid
)["uid"]
link_attachments([uid512, uid256, uid64], "award", award_uid)
now = datetime.now(timezone.utc).isoformat()
awards.update(
{
"uid": award_uid,
"attachment_uid_512": uid512,
"attachment_uid_256": uid256,
"attachment_uid_64": uid64,
"generated_at": now,
},
["uid"],
)
recompute_user_award_stats(receiver_uid)
users = get_table("users")
giver = users.find_one(uid=giver_uid) or {}
receiver = users.find_one(uid=receiver_uid) or {}
slug = row.get("slug", "")
create_notification(
receiver_uid,
"award",
f"@{giver.get('username', 'someone')} gave you an award",
giver_uid,
f"/profile/{receiver.get('username', '')}?tab=awards#award-{slug}",
)
audit.record_system(
"award.complete",
actor_kind="service",
target_type="award",
target_uid=award_uid,
target_label=description[:80],
summary=f"Award ready for {receiver.get('username', receiver_uid)}",
links=[
audit.target("award", award_uid, slug),
audit.target("user", receiver_uid, receiver.get("username")),
audit.target("user", giver_uid, giver.get("username")),
audit.job(job.get("uid", "")),
],
)
return {"award_uid": award_uid, "slug": slug}
def _fetch_image(self, api_key: str, description: str, totals: dict) -> bytes:
prompt = get_setting("award_image_prompt", AWARD_IMAGE_PROMPT_DEFAULT)
model = get_setting("award_image_model", AWARD_IMAGE_MODEL_DEFAULT)
size = AWARD_IMAGE_SIZE_DEFAULT
body = {
"model": model,
"prompt": f"{prompt}\n\n{description}".strip(),
"size": size,
"background": "transparent",
"output_format": "png",
}
url = f"{INTERNAL_BASE_URL}/openai/v1/images/generations"
headers = {"Content-Type": "application/json"}
if api_key:
headers["Authorization"] = f"Bearer {api_key}"
with stealth.stealth_sync_client(timeout=AWARD_GENERATION_TIMEOUT_SECONDS) as client:
response = client.post(url, json=body, headers=headers)
accumulate_usage(totals, response)
response.raise_for_status()
data = response.json()
items = data.get("data") or []
if not items:
raise RuntimeError("image API returned no data")
encoded = items[0].get("b64_json") or ""
if not encoded:
raise RuntimeError("image API returned no b64_json")
return base64.b64decode(encoded)
def _audit_failed(self, job: dict, error: str) -> None:
from devplacepy.services.audit import record as audit
payload = job.get("payload") or {}
award_uid = str(payload.get("award_uid", ""))
audit.record_system(
"award.failed",
actor_kind="service",
target_type="award",
target_uid=award_uid,
summary=f"Award generation failed: {error[:80]}",
metadata={"error": error[:500], "award_uid": award_uid},
result="failure",
links=[audit.job(job.get("uid", "")), audit.target("award", award_uid)],
)
def collect_metrics(self) -> dict:
base = super().collect_metrics()
base["stats"] = base.get("stats", []) + usage_metric_cards(get_award_usage())
return base
+17 -3
View File
@@ -52,6 +52,10 @@ The gateway records one row per upstream call (chat, vision, passthrough) and su
**Financial data is admin-only everywhere.** Any monetary figure (USD cost, pricing, spend, limit) is restricted to administrators; members and guests see only the percentage of quota used - this rule is enforced consistently across the profile card, `ai_correction`/`ai_modifier` usage displays, and the Devii cost tools.
## Image generation
`POST /openai/v1/images/generations` exposes an OpenAI-compatible image-generation endpoint. Clients send the generic model `molodetz-img-small` (`config.INTERNAL_IMAGE_MODEL`), which `handle_images` remaps to `gateway_image_model` exactly like chat remaps `molodetz` -> `gateway_model` (also remapped when `gateway_force_model` is on or the model is empty or `molodetz-img`). It defaults to OpenRouter's `black-forest-labs/flux-1.1-pro` at `https://openrouter.ai/api/v1/images/generations` (`config.IMAGE_*_DEFAULT`, $0.04 per image fallback). `handle_images` mirrors `handle_embeddings`: build the payload, forward via `_send`, and record one ledger row. The config fields are the **Images** group (`gateway_image_enabled` default on, `gateway_image_url`, `gateway_image_model`, `gateway_image_key`) plus the Pricing-group `gateway_image_price_per_call`. `effective_config()` falls the image key back to `gateway_api_key` then `OPENROUTER_API_KEY`. Usage is recorded with **`backend="image"`**; `usage.compute_cost` adds an `image` branch (flat per-call, native OpenRouter `cost` still preferred via `extract_image_usage`). `routing.image_overlay` resolves per-route provider/url/key and uses `price_input_per_m` as the per-image price. `routing.seed_default_image_routes()` (from `migrate_ai_gateway_settings`) idempotently seeds `molodetz-img-small` -> Flux on the `openrouter` provider when `OPENROUTER_API_KEY` is set. When `gateway_image_enabled` is off the endpoint returns 503 with no ledger row.
## Embeddings
`POST /openai/v1/embeddings` exposes an OpenAI-compatible text-embeddings model. Clients send the generic model `molodetz~embed` (`config.INTERNAL_EMBED_MODEL`), which `handle_embeddings` remaps to `gateway_embed_model` exactly like chat remaps `molodetz` -> `gateway_model` (also remapped when `gateway_force_model` is on or the model is empty). It defaults to OpenRouter's `qwen/qwen3-embedding-8b` at `https://openrouter.ai/api/v1/embeddings` (`config.EMBED_*_DEFAULT`, $0.01 per 1M input tokens). `handle_embeddings` mirrors `handle_chat` but is simpler: no vision augmentation and no streaming - build the payload, forward via `_send`, and record one ledger row through the same `finalize(...)` closure. The config fields are the **Embeddings** group (`gateway_embed_enabled` default on, `gateway_embed_url`, `gateway_embed_model`, `gateway_embed_key`) plus the Pricing-group `gateway_embed_price_input_per_m`. `effective_config()` falls the embed key back to `gateway_vision_key` then `OPENROUTER_API_KEY` (NOT `gateway_api_key`: that is the DeepSeek chat upstream key, whereas embeddings target OpenRouter like vision does). Usage is recorded with **`backend="embed"`**; `usage.compute_cost` adds an `embed` branch (input-only, completion always 0, native OpenRouter `cost` still preferred) and `Pricing` gained `embed_input_per_m`. `analytics.py` groups by `backend` generically, so embed rows roll up automatically; `caching_savings` counts only **non-native** chat rows (native-priced rows did not use the configured cache-hit/miss rates, so folding them in would report a fictional saving). When `gateway_embed_enabled` is off the endpoint returns 503 with no ledger row.
@@ -75,11 +79,11 @@ Layered ON TOP of the single-provider service config above, which stays THE impl
**Storage.** Two dataset tables, ensured in `init_db` via `routing.ensure_tables`, cross-worker cache-invalidated under the `"gateway_routing"` cache-version name (module-level `provider_store`/`model_store` over a shared `_ROUTING_CACHE`; writes `bump_cache_version` and clear the cache). They are admin config, NOT in `SOFT_DELETE_TABLES` - hard CRUD, mirroring `site_settings`:
- `gateway_providers` - named upstreams: `name` + `base_url` (chat-completions URL) + `api_key` + `is_active`; the embeddings URL is derived by swapping `/chat/completions` -> `/embeddings`.
- `gateway_models` - source->target routes: `source_model` (unique, what clients request) -> `provider` (blank = default) + `target_model` + `kind` (chat|embed) + optional `vision_provider`/`vision_model` for the **text+vision merge** (when set, image content is described by that vision model before forwarding) + `context_window` + its **own economy** (`price_cache_hit_per_m`/`price_cache_miss_per_m`/`price_output_per_m`/`price_input_per_m`, USD per 1M tokens) + `is_active`.
- `gateway_models` - source->target routes: `source_model` (unique, what clients request) -> `provider` (blank = default) + `target_model` + `kind` (chat|embed|image) + optional `vision_provider`/`vision_model` for the **text+vision merge** (when set, image content is described by that vision model before forwarding) + `context_window` + its **own economy** (`price_cache_hit_per_m`/`price_cache_miss_per_m`/`price_output_per_m`/`price_input_per_m`, USD per 1M tokens for chat/embed/vision; for image routes `price_input_per_m` is a flat USD per image, plus the tiered/off-peak fields below) + `is_active`.
**Resolution is a per-request overlay, not a fork.** At request time `routing.chat_overlay(requested, cfg)` / `routing.embed_overlay(requested, cfg)` resolve an active route by the requested model name and return a per-request OVERLAY dict of `gateway_*` cfg keys (`gateway_force_model`+`gateway_model`=target, `gateway_upstream_url`/`gateway_api_key` from the provider, the price keys, an augmented `gateway_model_context_map`, and vision overlay keys); `handle_chat`/`handle_embeddings` merge it onto the base cfg (`cfg = {**cfg, **overlay}`) BEFORE everything else, so the existing model-selection / `pricing_from_cfg` / `parse_context_map` / vision / url+key paths transparently use the route's provider, target model, pricing, vision model and context window. `_ensure` (the httpx pool / semaphore / breaker / vision cache) reads only the non-overlaid pool keys, so the connection pool is never churned per request.
**Resolution is a per-request overlay, not a fork.** At request time `routing.chat_overlay(requested, cfg)` / `routing.embed_overlay(requested, cfg)` / `routing.image_overlay(requested, cfg)` resolve an active route by the requested model name and return a per-request OVERLAY dict of `gateway_*` cfg keys (`gateway_force_model`+`gateway_model`=target, `gateway_upstream_url`/`gateway_api_key` from the provider, the price keys, an augmented `gateway_model_context_map`, and vision overlay keys); `handle_chat`/`handle_embeddings` merge it onto the base cfg (`cfg = {**cfg, **overlay}`) BEFORE everything else, so the existing model-selection / `pricing_from_cfg` / `parse_context_map` / vision / url+key paths transparently use the route's provider, target model, pricing, vision model and context window. `_ensure` (the httpx pool / semaphore / breaker / vision cache) reads only the non-overlaid pool keys, so the connection pool is never churned per request.
**No matching route = `None` overlay = byte-identical legacy behavior** - this is the "nobody feels the transformation" guarantee: `molodetz`/`molodetz~embed` and every existing caller are byte-identical when no route matches. A route with a blank provider overlays only model+pricing(+vision), keeping the default upstream url/key.
**No matching route = `None` overlay = byte-identical legacy behavior** - this is the "nobody feels the transformation" guarantee: `molodetz`/`molodetz~embed`/`molodetz-img-small` and every existing caller are byte-identical when no route matches. A route with a blank provider overlays only model+pricing(+vision), keeping the default upstream url/key.
**CRUD.** Admin JSON at `/admin/gateway/{providers,models}` (`routers/admin/gateway_configs.py`, `require_admin`, Pydantic `ProviderIn`/`ModelRouteIn` validation, accepts both JSON from `static/js/GatewayAdmin.js` and form from Devii), audited under `gateway.provider.*`/`gateway.model.*` (category `ai`), included in the admin package with the page at `/admin/gateway` (`templates/admin_gateway.html`, sidebar link, `admin_section="gateway"`).
@@ -88,3 +92,13 @@ Layered ON TOP of the single-provider service config above, which stays THE impl
When adding a routed value, overlay it as the matching `gateway_*` cfg key so the runtime needs no new branch.
**Tests:** `tests/unit/services/openai_gateway/routing.py` (overlay/economy/kind isolation), `tests/unit/services/openai_gateway/gateway.py::test_model_route_overrides_upstream` (end-to-end through `handle_chat`), and `tests/api/admin/gateway/` (admin CRUD + validation + role gating).
## Tiered (context-length) and off-peak pricing (model-agnostic variable pricing)
Real providers sometimes charge more than a flat per-1M rate for one component: a rate that jumps once a request crosses a context-length threshold, or a fixed time-of-day discount window. This is layered on top of the cache-hit/cache-miss/output shape (which is already exactly DeepSeek's real billing model - see below), on the SAME `gateway_models` route row, so it stays fully provider-agnostic and opt-in per route.
- **Per-route fields** (each optional/zero by default = feature off, so an unconfigured route is byte-identical to before): `context_tier_threshold_tokens` (0 disables tiering; when the request's input token count exceeds it, `price_*_per_m_tier2` rates apply instead of the tier-1 rates above) and the four nullable `price_cache_hit_per_m_tier2`/`price_cache_miss_per_m_tier2`/`price_output_per_m_tier2`/`price_input_per_m_tier2` (a component left `None` keeps its tier-1 rate even above threshold - so a provider that only re-prices input above a size threshold, keeping output flat, needs just one tier2 field set), plus `off_peak_start_minute`/`off_peak_end_minute` (nullable, UTC minutes-since-midnight, both-or-neither enforced by `ModelRouteIn`'s model validator; a window where start > end wraps past midnight, e.g. DeepSeek's old V3/R1-era 16:30-00:30 UTC window) and `off_peak_discount_pct` (0-100, multiplies whichever tier rate is active). The same fields double for vision (`price_input/output_per_m_tier2`) and embed (`price_input_per_m_tier2`) exactly like their tier-1 counterparts already do.
- **Overlay + computation.** `chat_overlay`/`embed_overlay` propagate all of these into the per-request cfg dict under `gateway_*` keys exactly like the existing price fields (section above); `usage.pricing_from_cfg` reads them generically (defaulting to `None`/`0`/disabled when absent), so **Layer A (the single global flat Pricing config fields) never gains these dimensions** - only a `gateway_models` route can enable them, preserving the "no matching route = byte-identical legacy behavior" guarantee. `usage.compute_cost` selects tier1 vs tier2 per rate component (`_tiered_rate`) based on whether `norm["prompt"]` (billable input tokens) exceeds the threshold, then applies the off-peak discount (`_effective_rate`/`_off_peak_active`, UTC wraparound-aware) to whichever rate was selected. **The response header format and `compute_cost`'s return shape (`total, input_cost, output_cost, native`) are unchanged** - this is purely an internal rate-selection step before the existing input/output split math runs; a native upstream `cost` (OpenRouter) still overrides the modeled total exactly as before.
- **Migration.** New `gateway_models` columns are added via `has_column`/`create_column_by_example` in `routing.ensure_tables()` (the `CREATE TABLE IF NOT EXISTS` DDL string alone would never reach a pre-existing table - see the `database/CLAUDE.md` column-ensure idiom).
- **Admin UI.** `/admin/gateway`'s model-route form has a "Tiered / off-peak pricing (optional)" subsection; off-peak start/end render as `<input type="time">` (converted to/from UTC minutes-of-day by `GatewayAdmin.js`), and the routes table shows `tiered`/`off-peak` badges when a route has either dimension configured.
- **DeepSeek's real pricing is already the tier-1 shape, not a new dimension.** DeepSeek's actual API (verified against `api-docs.deepseek.com/quick_start/pricing`) bills three flat per-1M rates - cache-hit input, cache-miss input, output - with no current context-length tier or off-peak window for the V4 models; that shape was already fully modeled by the pre-existing `chat_cache_hit_per_m`/`chat_cache_miss_per_m`/`chat_output_per_m` fields before this section's tier2/off-peak fields existed. `routing.seed_default_deepseek_routes()` (called once from `database.migrate_ai_gateway_settings()` at the end of `init_db()`) idempotently inserts two ready-made routes - `deepseek-v4-flash` (`$0.0028`/`$0.14`/`$0.28` per 1M, 1M context) and `deepseek-v4-pro` (`$0.003625`/`$0.435`/`$0.87` per 1M, 1M context) - only when that `source_model` row does not already exist, so a caller or Devii can request either name explicitly and get correctly-priced, decoupled from whatever the single global `gateway_model` default happens to be set to (switching that global setting between the two real models does NOT retroactively fix the flat Pricing config fields - the seeded routes are the model-agnostic, always-correct way to reference a specific priced model). Neither seeded route sets the tier2/off-peak fields (DeepSeek does not use them today); an admin can add them later on the same row if DeepSeek (or any other provider routed here) introduces such pricing.
@@ -15,6 +15,10 @@ VISION_CACHE_SIZE_DEFAULT = 256
EMBED_URL_DEFAULT = "https://openrouter.ai/api/v1/embeddings"
EMBED_MODEL_DEFAULT = "qwen/qwen3-embedding-8b"
IMAGE_URL_DEFAULT = "https://openrouter.ai/api/v1/images"
IMAGE_MODEL_DEFAULT = "black-forest-labs/flux.2-pro"
IMAGE_PRICE_PER_CALL_DEFAULT = 0.04
VISION_INSTRUCTION = (
"Describe this image in detail. Note objects, people, scene, any visible "
"text, layout, colors, and anything else that could be relevant for "
+173 -1
View File
@@ -13,11 +13,16 @@ from fastapi.responses import JSONResponse, Response, StreamingResponse
from devplacepy import stealth
from devplacepy.services.openai_gateway import config
from devplacepy.services.openai_gateway.reliability import CircuitBreaker, retry_send
from devplacepy.services.openai_gateway.routing import chat_overlay, embed_overlay
from devplacepy.services.openai_gateway.routing import (
chat_overlay,
embed_overlay,
image_overlay,
)
from devplacepy.services.openai_gateway.system_message import apply_system_directives
from devplacepy.services.openai_gateway.usage import (
GatewayUsageLedger,
classify_error,
extract_image_usage,
extract_params,
parse_context_map,
pricing_from_cfg,
@@ -106,6 +111,7 @@ class GatewayRuntime:
self.peak_in_flight = 0
self.vision_calls = 0
self.embed_calls = 0
self.image_calls = 0
self.last_status = 0
self.last_latency_ms = 0
@@ -537,6 +543,171 @@ class GatewayRuntime:
log(f"embed POST -> 200 ({timing['upstream_latency_ms']:.0f}ms)")
return JSONResponse(content=data, headers=resp_headers)
async def handle_images(
self, body: dict, cfg: dict, owner: tuple, user_agent: str, log=None
):
log = log or (lambda message: None)
overlay = image_overlay(body.get("model"), cfg)
if overlay:
cfg = {**cfg, **overlay}
log(
f"routed image model {body.get('model')!r} -> {cfg['gateway_image_model']!r}"
)
if not cfg["gateway_image_enabled"]:
from devplacepy.services.audit import record as audit
from devplacepy.services.openai_gateway.usage import audit_actor_for
actor_kind, actor_uid, actor_role = audit_actor_for(owner[0], owner[1])
audit.record_system(
"ai.gateway.call",
actor_kind=actor_kind,
actor_uid=actor_uid,
actor_role=actor_role,
origin="api",
result="denied",
summary="image generation disabled",
metadata={
"backend": "image",
"endpoint": "images/generations",
"owner_kind": owner[0],
"owner_id": owner[1],
},
)
return JSONResponse(
status_code=503,
content={
"error": {
"message": "Image generation is disabled",
"type": "images_disabled",
}
},
)
client, sem = self._ensure(cfg)
pricing = pricing_from_cfg(cfg)
context_map = parse_context_map(cfg.get("gateway_model_context_map"))
params = extract_params(body)
handle_start = time.monotonic()
requested = body.get("model")
if (
cfg["gateway_force_model"]
or not requested
or requested in ("molodetz-img-small", "molodetz-img")
):
model = cfg["gateway_image_model"]
else:
model = requested
payload = dict(body)
payload["model"] = model
headers = {"Content-Type": "application/json"}
if cfg["gateway_image_key"]:
headers["Authorization"] = f"Bearer {cfg['gateway_image_key']}"
else:
log(
"No upstream image API key configured (gateway_image_key / gateway_api_key / OPENROUTER_API_KEY); upstream will likely reject the request"
)
resp, exc, timing = await self._send(
client,
sem,
"POST",
cfg["gateway_image_url"],
headers,
cfg,
log,
json_body=payload,
)
base = {
"owner_kind": owner[0],
"owner_id": owner[1],
"backend": "image",
"endpoint": "images/generations",
"model": model,
"user_agent": user_agent,
**params,
**timing,
}
def finalize(status_code, success, category, usage=None):
base["total_latency_ms"] = round(
(time.monotonic() - handle_start) * 1000, 3
)
base["gateway_overhead_ms"] = round(
max(
base["total_latency_ms"]
- timing["upstream_latency_ms"]
- timing["queue_wait_ms"],
0.0,
),
3,
)
base["status_code"] = status_code
base["success"] = success
base["error_category"] = category
base["usage"] = usage
row = self._ledger.record(base, pricing, context_map)
return usage_response_headers(row)
if timing["circuit_open"]:
resp_headers = finalize(503, False, "circuit_open")
return JSONResponse(
status_code=503,
content={
"error": {
"message": "Upstream temporarily unavailable",
"type": "circuit_open",
}
},
headers=resp_headers,
)
if exc is not None:
resp_headers = finalize(502, False, classify_error(0, exc))
return JSONResponse(
status_code=502,
content={
"error": {
"message": f"Upstream connection failed: {exc}",
"type": "upstream_error",
}
},
headers=resp_headers,
)
if resp.status_code != 200:
resp_headers = finalize(
resp.status_code,
False,
classify_error(resp.status_code, None, resp.text),
)
log(f"image upstream POST -> {resp.status_code}: {resp.text[:300]}")
return JSONResponse(
status_code=resp.status_code,
content={"error": {"message": resp.text, "type": "upstream_error"}},
headers=resp_headers,
)
try:
data = resp.json()
except ValueError:
self.errors += 1
resp_headers = finalize(502, False, "gateway")
log("image upstream returned 200 but body was not valid JSON")
return JSONResponse(
status_code=502,
content={
"error": {
"message": "invalid upstream response",
"type": "upstream_error",
}
},
headers=resp_headers,
)
self.image_calls += 1
resp_headers = finalize(200, True, None, extract_image_usage(data))
log(f"image POST -> 200 ({timing['upstream_latency_ms']:.0f}ms)")
return JSONResponse(content=data, headers=resp_headers)
async def handle_passthrough(
self,
method: str,
@@ -661,6 +832,7 @@ class GatewayRuntime:
"peak_in_flight": self.peak_in_flight,
"vision_calls": self.vision_calls,
"embed_calls": self.embed_calls,
"image_calls": self.image_calls,
"last_status": self.last_status,
"last_latency_ms": self.last_latency_ms,
"pool": self._instances,
+288 -3
View File
@@ -7,7 +7,7 @@ from dataclasses import dataclass
from datetime import datetime, timezone
from typing import Optional
from pydantic import BaseModel, Field, field_validator
from pydantic import BaseModel, Field, field_validator, model_validator
from devplacepy.database import (
bump_cache_version,
@@ -21,7 +21,7 @@ logger = logging.getLogger(__name__)
PROVIDERS_TABLE = "gateway_providers"
MODELS_TABLE = "gateway_models"
CACHE_NAME = "gateway_routing"
KINDS = ("chat", "embed")
KINDS = ("chat", "embed", "image")
_ROUTING_CACHE: dict = {}
@@ -47,6 +47,27 @@ def _embed_url_from_base(base_url: str) -> str:
return base_url
def _image_url_from_base(base_url: str) -> str:
base_url = (base_url or "").strip()
if not base_url:
return ""
if base_url.endswith("/chat/completions"):
return base_url[: -len("/chat/completions")] + "/images"
return base_url
MODEL_TIER2_COLUMNS = (
"context_tier_threshold_tokens",
"price_cache_hit_per_m_tier2",
"price_cache_miss_per_m_tier2",
"price_output_per_m_tier2",
"price_input_per_m_tier2",
"off_peak_start_minute",
"off_peak_end_minute",
"off_peak_discount_pct",
)
def ensure_tables() -> None:
db.query(
"CREATE TABLE IF NOT EXISTS "
@@ -62,8 +83,17 @@ def ensure_tables() -> None:
"vision_model TEXT, context_window INTEGER DEFAULT 0, "
"price_cache_hit_per_m REAL DEFAULT 0, price_cache_miss_per_m REAL DEFAULT 0, "
"price_output_per_m REAL DEFAULT 0, price_input_per_m REAL DEFAULT 0, "
"context_tier_threshold_tokens INTEGER DEFAULT 0, "
"price_cache_hit_per_m_tier2 REAL, price_cache_miss_per_m_tier2 REAL, "
"price_output_per_m_tier2 REAL, price_input_per_m_tier2 REAL, "
"off_peak_start_minute INTEGER, off_peak_end_minute INTEGER, "
"off_peak_discount_pct REAL DEFAULT 0, "
"is_active INTEGER DEFAULT 1, created_at TEXT, updated_at TEXT)"
)
models_table = get_table(MODELS_TABLE)
for column in MODEL_TIER2_COLUMNS:
if not models_table.has_column(column):
models_table.create_column_by_example(column, 0.0)
try:
db.query(
"CREATE UNIQUE INDEX IF NOT EXISTS idx_gateway_providers_name ON "
@@ -121,6 +151,14 @@ class ModelRouteIn(BaseModel):
price_cache_miss_per_m: float = Field(default=0.0, ge=0)
price_output_per_m: float = Field(default=0.0, ge=0)
price_input_per_m: float = Field(default=0.0, ge=0)
context_tier_threshold_tokens: int = Field(default=0, ge=0, le=100_000_000)
price_cache_hit_per_m_tier2: Optional[float] = Field(default=None, ge=0)
price_cache_miss_per_m_tier2: Optional[float] = Field(default=None, ge=0)
price_output_per_m_tier2: Optional[float] = Field(default=None, ge=0)
price_input_per_m_tier2: Optional[float] = Field(default=None, ge=0)
off_peak_start_minute: Optional[int] = Field(default=None, ge=0, le=1439)
off_peak_end_minute: Optional[int] = Field(default=None, ge=0, le=1439)
off_peak_discount_pct: float = Field(default=0.0, ge=0, le=100)
is_active: bool = True
@field_validator("source_model", "target_model")
@@ -141,9 +179,19 @@ class ModelRouteIn(BaseModel):
def _clean_kind(cls, value: str) -> str:
value = (value or "chat").strip().lower()
if value not in KINDS:
raise ValueError("Kind must be 'chat' or 'embed'")
raise ValueError("Kind must be 'chat', 'embed', or 'image'")
return value
@model_validator(mode="after")
def _check_off_peak_window(self) -> "ModelRouteIn":
has_start = self.off_peak_start_minute is not None
has_end = self.off_peak_end_minute is not None
if has_start != has_end:
raise ValueError(
"Off-peak start and end minute must both be set, or both left blank"
)
return self
@dataclass(frozen=True)
class ModelRoute:
@@ -158,9 +206,27 @@ class ModelRoute:
price_cache_miss_per_m: float
price_output_per_m: float
price_input_per_m: float
context_tier_threshold_tokens: int
price_cache_hit_per_m_tier2: Optional[float]
price_cache_miss_per_m_tier2: Optional[float]
price_output_per_m_tier2: Optional[float]
price_input_per_m_tier2: Optional[float]
off_peak_start_minute: Optional[int]
off_peak_end_minute: Optional[int]
off_peak_discount_pct: float
is_active: bool
def _opt_float(row: dict, key: str) -> Optional[float]:
value = row.get(key)
return float(value) if value is not None else None
def _opt_int(row: dict, key: str) -> Optional[int]:
value = row.get(key)
return int(value) if value is not None else None
def _route_from_row(row: dict) -> ModelRoute:
return ModelRoute(
source_model=str(row.get("source_model") or ""),
@@ -174,6 +240,16 @@ def _route_from_row(row: dict) -> ModelRoute:
price_cache_miss_per_m=float(row.get("price_cache_miss_per_m") or 0.0),
price_output_per_m=float(row.get("price_output_per_m") or 0.0),
price_input_per_m=float(row.get("price_input_per_m") or 0.0),
context_tier_threshold_tokens=int(
row.get("context_tier_threshold_tokens") or 0
),
price_cache_hit_per_m_tier2=_opt_float(row, "price_cache_hit_per_m_tier2"),
price_cache_miss_per_m_tier2=_opt_float(row, "price_cache_miss_per_m_tier2"),
price_output_per_m_tier2=_opt_float(row, "price_output_per_m_tier2"),
price_input_per_m_tier2=_opt_float(row, "price_input_per_m_tier2"),
off_peak_start_minute=_opt_int(row, "off_peak_start_minute"),
off_peak_end_minute=_opt_int(row, "off_peak_end_minute"),
off_peak_discount_pct=float(row.get("off_peak_discount_pct") or 0.0),
is_active=_as_bool(row.get("is_active", 1)),
)
@@ -291,6 +367,14 @@ class ModelStore:
"price_cache_miss_per_m": payload.price_cache_miss_per_m,
"price_output_per_m": payload.price_output_per_m,
"price_input_per_m": payload.price_input_per_m,
"context_tier_threshold_tokens": payload.context_tier_threshold_tokens,
"price_cache_hit_per_m_tier2": payload.price_cache_hit_per_m_tier2,
"price_cache_miss_per_m_tier2": payload.price_cache_miss_per_m_tier2,
"price_output_per_m_tier2": payload.price_output_per_m_tier2,
"price_input_per_m_tier2": payload.price_input_per_m_tier2,
"off_peak_start_minute": payload.off_peak_start_minute,
"off_peak_end_minute": payload.off_peak_end_minute,
"off_peak_discount_pct": payload.off_peak_discount_pct,
"is_active": 1 if payload.is_active else 0,
"updated_at": _now(),
}
@@ -318,6 +402,63 @@ class ModelStore:
provider_store = ProviderStore()
model_store = ModelStore()
DEEPSEEK_DEFAULT_ROUTES = (
{
"source_model": "deepseek-v4-flash",
"target_model": "deepseek-v4-flash",
"context_window": 1_048_576,
"price_cache_hit_per_m": 0.0028,
"price_cache_miss_per_m": 0.14,
"price_output_per_m": 0.28,
},
{
"source_model": "deepseek-v4-pro",
"target_model": "deepseek-v4-pro",
"context_window": 1_048_576,
"price_cache_hit_per_m": 0.003625,
"price_cache_miss_per_m": 0.435,
"price_output_per_m": 0.87,
},
)
def seed_default_deepseek_routes() -> None:
ensure_tables()
table = get_table(MODELS_TABLE)
seeded = False
for defaults in DEEPSEEK_DEFAULT_ROUTES:
if table.find_one(source_model=defaults["source_model"]):
continue
record = {
"source_model": defaults["source_model"],
"provider": "",
"target_model": defaults["target_model"],
"kind": "chat",
"vision_provider": "",
"vision_model": "",
"context_window": defaults["context_window"],
"price_cache_hit_per_m": defaults["price_cache_hit_per_m"],
"price_cache_miss_per_m": defaults["price_cache_miss_per_m"],
"price_output_per_m": defaults["price_output_per_m"],
"price_input_per_m": 0.0,
"context_tier_threshold_tokens": 0,
"price_cache_hit_per_m_tier2": None,
"price_cache_miss_per_m_tier2": None,
"price_output_per_m_tier2": None,
"price_input_per_m_tier2": None,
"off_peak_start_minute": None,
"off_peak_end_minute": None,
"off_peak_discount_pct": 0.0,
"is_active": 1,
"created_at": _now(),
"updated_at": _now(),
}
table.insert(record)
seeded = True
if seeded:
bump_cache_version(CACHE_NAME)
_ROUTING_CACHE.clear()
def _provider_overlay(name: str, base_key: str, url_key: str, overlay: dict) -> None:
provider = provider_store.get(name)
@@ -339,6 +480,13 @@ def chat_overlay(requested_model: Optional[str], base_cfg: dict) -> Optional[dic
"gateway_price_cache_hit_per_m": route.price_cache_hit_per_m,
"gateway_price_cache_miss_per_m": route.price_cache_miss_per_m,
"gateway_price_output_per_m": route.price_output_per_m,
"gateway_price_cache_hit_per_m_tier2": route.price_cache_hit_per_m_tier2,
"gateway_price_cache_miss_per_m_tier2": route.price_cache_miss_per_m_tier2,
"gateway_price_output_per_m_tier2": route.price_output_per_m_tier2,
"gateway_context_tier_threshold_tokens": route.context_tier_threshold_tokens,
"gateway_off_peak_start_minute": route.off_peak_start_minute,
"gateway_off_peak_end_minute": route.off_peak_end_minute,
"gateway_off_peak_discount_pct": route.off_peak_discount_pct,
}
if route.provider:
_provider_overlay(
@@ -355,6 +503,12 @@ def chat_overlay(requested_model: Optional[str], base_cfg: dict) -> Optional[dic
overlay["gateway_vision_model"] = route.vision_model
overlay["gateway_vision_price_input_per_m"] = route.price_input_per_m
overlay["gateway_vision_price_output_per_m"] = route.price_output_per_m
overlay["gateway_vision_price_input_per_m_tier2"] = (
route.price_input_per_m_tier2
)
overlay["gateway_vision_price_output_per_m_tier2"] = (
route.price_output_per_m_tier2
)
vision_provider = route.vision_provider or route.provider
if vision_provider:
_provider_overlay(
@@ -371,6 +525,11 @@ def embed_overlay(requested_model: Optional[str], base_cfg: dict) -> Optional[di
"gateway_force_model": True,
"gateway_embed_model": route.target_model,
"gateway_embed_price_input_per_m": route.price_input_per_m,
"gateway_embed_price_input_per_m_tier2": route.price_input_per_m_tier2,
"gateway_context_tier_threshold_tokens": route.context_tier_threshold_tokens,
"gateway_off_peak_start_minute": route.off_peak_start_minute,
"gateway_off_peak_end_minute": route.off_peak_end_minute,
"gateway_off_peak_discount_pct": route.off_peak_discount_pct,
}
provider = provider_store.get(route.provider) if route.provider else None
if provider:
@@ -379,3 +538,129 @@ def embed_overlay(requested_model: Optional[str], base_cfg: dict) -> Optional[di
if provider.get("api_key"):
overlay["gateway_embed_key"] = provider["api_key"]
return overlay
def image_overlay(requested_model: Optional[str], base_cfg: dict) -> Optional[dict]:
route = model_store.resolve(requested_model, "image")
if route is None:
return None
overlay: dict = {
"gateway_force_model": True,
"gateway_image_model": route.target_model,
"gateway_image_price_per_call": route.price_input_per_m,
"gateway_off_peak_start_minute": route.off_peak_start_minute,
"gateway_off_peak_end_minute": route.off_peak_end_minute,
"gateway_off_peak_discount_pct": route.off_peak_discount_pct,
}
provider = provider_store.get(route.provider) if route.provider else None
if provider:
if provider.get("base_url"):
overlay["gateway_image_url"] = _image_url_from_base(provider["base_url"])
if provider.get("api_key"):
overlay["gateway_image_key"] = provider["api_key"]
return overlay
OPENROUTER_PROVIDER_NAME = "openrouter"
FLUX_TARGET_MODEL = "black-forest-labs/flux.2-pro"
IMAGE_SOURCE_MODEL = "molodetz-img-small"
RETIRED_IMAGE_MODELS = {
"black-forest-labs/flux-1.1-pro": FLUX_TARGET_MODEL,
"black-forest-labs/flux-schnell": "black-forest-labs/flux.2-klein-4b",
}
LEGACY_IMAGE_URL = "https://openrouter.ai/api/v1/images/generations"
def seed_default_image_routes() -> None:
import os
from devplacepy.services.openai_gateway import config as gw_config
ensure_tables()
openrouter_key = os.environ.get("OPENROUTER_API_KEY", "")
if openrouter_key:
existing_provider = provider_store.get(OPENROUTER_PROVIDER_NAME)
if existing_provider is None:
provider_store.set(
ProviderIn(
name=OPENROUTER_PROVIDER_NAME,
base_url="https://openrouter.ai/api/v1/chat/completions",
api_key=openrouter_key,
is_active=True,
)
)
models_table = get_table(MODELS_TABLE)
if models_table.find_one(source_model=IMAGE_SOURCE_MODEL):
return
record = {
"source_model": IMAGE_SOURCE_MODEL,
"provider": OPENROUTER_PROVIDER_NAME if openrouter_key else "",
"target_model": FLUX_TARGET_MODEL,
"kind": "image",
"vision_provider": "",
"vision_model": "",
"context_window": 0,
"price_cache_hit_per_m": 0.0,
"price_cache_miss_per_m": 0.0,
"price_output_per_m": 0.0,
"price_input_per_m": gw_config.IMAGE_PRICE_PER_CALL_DEFAULT,
"context_tier_threshold_tokens": 0,
"price_cache_hit_per_m_tier2": None,
"price_cache_miss_per_m_tier2": None,
"price_output_per_m_tier2": None,
"price_input_per_m_tier2": None,
"off_peak_start_minute": None,
"off_peak_end_minute": None,
"off_peak_discount_pct": 0.0,
"is_active": 1,
"created_at": _now(),
"updated_at": _now(),
}
models_table.insert(record)
bump_cache_version(CACHE_NAME)
_ROUTING_CACHE.clear()
def migrate_retired_image_gateway() -> None:
from devplacepy.database import get_setting, set_setting
from devplacepy.services.openai_gateway import config as gw_config
ensure_tables()
models_table = get_table(MODELS_TABLE)
now = _now()
changed = False
for row in models_table.find(kind="image"):
target = str(row.get("target_model") or "")
replacement = RETIRED_IMAGE_MODELS.get(target)
if not replacement:
continue
models_table.update(
{
"source_model": row["source_model"],
"target_model": replacement,
"updated_at": now,
},
["source_model"],
)
changed = True
image_url = get_setting("gateway_image_url", "")
if image_url in (LEGACY_IMAGE_URL, ""):
set_setting("gateway_image_url", gw_config.IMAGE_URL_DEFAULT)
changed = True
elif image_url.endswith("/images/generations"):
set_setting(
"gateway_image_url",
image_url[: -len("/generations")],
)
changed = True
image_model = get_setting("gateway_image_model", "")
replacement = RETIRED_IMAGE_MODELS.get(image_model)
if replacement:
set_setting("gateway_image_model", replacement)
changed = True
elif not image_model:
set_setting("gateway_image_model", gw_config.IMAGE_MODEL_DEFAULT)
changed = True
if changed:
bump_cache_version(CACHE_NAME)
_ROUTING_CACHE.clear()
@@ -172,6 +172,38 @@ class GatewayService(BaseService):
help="The key currently in use; falls back to the vision/OPENROUTER key on boot (embeddings default to OpenRouter). Editable.",
group="Embeddings",
),
ConfigField(
"gateway_image_enabled",
"Image generation",
type="bool",
default=True,
help="Expose image generation at /openai/v1/images/generations.",
group="Images",
),
ConfigField(
"gateway_image_url",
"Images URL",
type="url",
default=config.IMAGE_URL_DEFAULT,
help="OpenAI-compatible images/generations endpoint requests are forwarded to.",
group="Images",
),
ConfigField(
"gateway_image_model",
"Images model",
type="str",
default=config.IMAGE_MODEL_DEFAULT,
help="Image model sent upstream. Clients request it as molodetz-img-small.",
group="Images",
),
ConfigField(
"gateway_image_key",
"Images API key",
type="str",
default="",
help="The key currently in use; falls back to gateway_api_key / OPENROUTER_API_KEY on boot. Editable.",
group="Images",
),
ConfigField(
"gateway_require_auth",
"Require authentication",
@@ -267,6 +299,15 @@ class GatewayService(BaseService):
help="Fallback only; used when the embeddings upstream returns no native cost.",
group="Pricing",
),
ConfigField(
"gateway_image_price_per_call",
"Images price / call ($)",
type="float",
default=config.IMAGE_PRICE_PER_CALL_DEFAULT,
minimum=0,
help="Fallback only; used when the image upstream returns no native cost.",
group="Pricing",
),
ConfigField(
"gateway_rsearch_cost_per_call",
"rsearch cost / call ($)",
@@ -355,6 +396,17 @@ class GatewayService(BaseService):
or cfg["gateway_vision_key"]
or os.environ.get("OPENROUTER_API_KEY", "")
)
cfg["gateway_image_key"] = (
cfg["gateway_image_key"]
or cfg["gateway_api_key"]
or os.environ.get("OPENROUTER_API_KEY", "")
)
if not cfg.get("gateway_image_url"):
upstream = cfg.get("gateway_upstream_url", "")
if upstream.endswith("/chat/completions"):
cfg["gateway_image_url"] = (
upstream[: -len("/chat/completions")] + "/images"
)
return cfg
def authorize(self, request: Request) -> bool:
@@ -429,6 +481,18 @@ class GatewayService(BaseService):
return await runtime.handle_embeddings(
body, cfg, owner, user_agent, self.log
)
if subpath == "images/generations" and request.method == "POST":
try:
body = await request.json()
except Exception:
self.log("Rejected images request: invalid JSON body")
raise HTTPException(status_code=400, detail="Invalid JSON body")
if not isinstance(body, dict):
self.log("Rejected images request: JSON body was not an object")
raise HTTPException(status_code=400, detail="Invalid JSON body")
return await runtime.handle_images(
body, cfg, owner, user_agent, self.log
)
body = await request.body()
content_type = request.headers.get("content-type", "")
return await runtime.handle_passthrough(
@@ -473,6 +537,7 @@ class GatewayService(BaseService):
"peak_in_flight": 0,
"vision_calls": 0,
"embed_calls": 0,
"image_calls": 0,
"last_status": 0,
"last_latency_ms": 0,
"pool": 0,
@@ -485,12 +550,14 @@ class GatewayService(BaseService):
{"label": "In flight", "value": m["in_flight"]},
{"label": "Vision calls", "value": m["vision_calls"]},
{"label": "Embed calls", "value": m["embed_calls"]},
{"label": "Image calls", "value": m["image_calls"]},
{"label": "Last status", "value": m["last_status"] or "-"},
{"label": "Last latency", "value": f"{m['last_latency_ms']} ms"},
{"label": "Pool size", "value": m["pool"]},
{"label": "Circuit", "value": "open" if m["circuit_open"] else "closed"},
{"label": "Model", "value": cfg["gateway_model"]},
{"label": "Embed model", "value": cfg["gateway_embed_model"]},
{"label": "Image model", "value": cfg["gateway_image_model"]},
{"label": "Requests 24h", "value": s["requests"]},
{"label": "Success 24h", "value": f"{s['success_pct']}%"},
{"label": "Error rate 24h", "value": f"{s['error_pct']}%"},
+160 -8
View File
@@ -36,6 +36,31 @@ class Pricing:
vision_input_per_m: float
vision_output_per_m: float
embed_input_per_m: float
image_per_call: float = 0.0
chat_cache_hit_per_m_tier2: Optional[float] = None
chat_cache_miss_per_m_tier2: Optional[float] = None
chat_output_per_m_tier2: Optional[float] = None
vision_input_per_m_tier2: Optional[float] = None
vision_output_per_m_tier2: Optional[float] = None
embed_input_per_m_tier2: Optional[float] = None
context_tier_threshold_tokens: int = 0
off_peak_start_minute: Optional[int] = None
off_peak_end_minute: Optional[int] = None
off_peak_discount_pct: float = 0.0
def _cfg_float(cfg: dict, key: str) -> Optional[float]:
value = cfg.get(key)
if isinstance(value, (int, float)) and not isinstance(value, bool):
return float(value)
return None
def _cfg_int(cfg: dict, key: str) -> Optional[int]:
value = cfg.get(key)
if isinstance(value, (int, float)) and not isinstance(value, bool):
return int(value)
return None
def pricing_from_cfg(cfg: dict) -> Pricing:
@@ -71,6 +96,34 @@ def pricing_from_cfg(cfg: dict) -> Pricing:
config.EMBED_PRICE_INPUT_PER_M_DEFAULT,
)
),
image_per_call=float(
cfg.get(
"gateway_image_price_per_call",
config.IMAGE_PRICE_PER_CALL_DEFAULT,
)
),
chat_cache_hit_per_m_tier2=_cfg_float(
cfg, "gateway_price_cache_hit_per_m_tier2"
),
chat_cache_miss_per_m_tier2=_cfg_float(
cfg, "gateway_price_cache_miss_per_m_tier2"
),
chat_output_per_m_tier2=_cfg_float(cfg, "gateway_price_output_per_m_tier2"),
vision_input_per_m_tier2=_cfg_float(
cfg, "gateway_vision_price_input_per_m_tier2"
),
vision_output_per_m_tier2=_cfg_float(
cfg, "gateway_vision_price_output_per_m_tier2"
),
embed_input_per_m_tier2=_cfg_float(
cfg, "gateway_embed_price_input_per_m_tier2"
),
context_tier_threshold_tokens=int(
cfg.get("gateway_context_tier_threshold_tokens", 0) or 0
),
off_peak_start_minute=_cfg_int(cfg, "gateway_off_peak_start_minute"),
off_peak_end_minute=_cfg_int(cfg, "gateway_off_peak_end_minute"),
off_peak_discount_pct=float(cfg.get("gateway_off_peak_discount_pct", 0.0) or 0.0),
)
@@ -88,6 +141,18 @@ def parse_context_map(raw: Any) -> dict[str, int]:
return dict(config.MODEL_CONTEXT_MAP_DEFAULT)
def extract_image_usage(data: Optional[dict]) -> dict:
data = data or {}
usage = data.get("usage") if isinstance(data.get("usage"), dict) else {}
cost = usage.get("cost")
if cost is None:
cost = data.get("cost")
result: dict = {}
if isinstance(cost, (int, float)) and not isinstance(cost, bool):
result["cost"] = float(cost)
return result
def normalize_usage(usage: Optional[dict]) -> dict:
usage = usage or {}
prompt = int(usage.get("prompt_tokens", 0) or 0)
@@ -117,21 +182,108 @@ def normalize_usage(usage: Optional[dict]) -> dict:
}
def _off_peak_active(pricing: Pricing, now: Optional[datetime] = None) -> bool:
if pricing.off_peak_start_minute is None or pricing.off_peak_end_minute is None:
return False
moment = now or _now()
minute_of_day = moment.hour * 60 + moment.minute
start, end = pricing.off_peak_start_minute, pricing.off_peak_end_minute
if start == end:
return True
if start < end:
return start <= minute_of_day < end
return minute_of_day >= start or minute_of_day < end
def _tiered_rate(base: float, tier2: Optional[float], use_tier2: bool) -> float:
if use_tier2 and tier2 is not None:
return tier2
return base
def _effective_rate(
base: float,
tier2: Optional[float],
use_tier2: bool,
pricing: Pricing,
off_peak: bool,
) -> float:
rate = _tiered_rate(base, tier2, use_tier2)
if off_peak and pricing.off_peak_discount_pct > 0:
rate = rate * (1 - min(pricing.off_peak_discount_pct, 100.0) / 100.0)
return rate
def compute_cost(
usage: dict, norm: dict, pricing: Pricing, backend: str
usage: dict,
norm: dict,
pricing: Pricing,
backend: str,
now: Optional[datetime] = None,
) -> tuple[float, float, float, bool]:
threshold = pricing.context_tier_threshold_tokens
use_tier2 = bool(threshold > 0 and norm["prompt"] > threshold)
off_peak = _off_peak_active(pricing, now)
if backend == "vision":
input_cost = norm["prompt"] / PER_MILLION * pricing.vision_input_per_m
output_cost = norm["completion"] / PER_MILLION * pricing.vision_output_per_m
input_rate = _effective_rate(
pricing.vision_input_per_m,
pricing.vision_input_per_m_tier2,
use_tier2,
pricing,
off_peak,
)
output_rate = _effective_rate(
pricing.vision_output_per_m,
pricing.vision_output_per_m_tier2,
use_tier2,
pricing,
off_peak,
)
input_cost = norm["prompt"] / PER_MILLION * input_rate
output_cost = norm["completion"] / PER_MILLION * output_rate
elif backend == "embed":
input_cost = norm["prompt"] / PER_MILLION * pricing.embed_input_per_m
input_rate = _effective_rate(
pricing.embed_input_per_m,
pricing.embed_input_per_m_tier2,
use_tier2,
pricing,
off_peak,
)
input_cost = norm["prompt"] / PER_MILLION * input_rate
output_cost = 0.0
elif backend == "image":
rate = pricing.image_per_call
if off_peak and pricing.off_peak_discount_pct > 0:
rate = rate * (1 - min(pricing.off_peak_discount_pct, 100.0) / 100.0)
input_cost = rate
output_cost = 0.0
else:
input_cost = (
norm["cache_hit"] / PER_MILLION * pricing.chat_cache_hit_per_m
+ norm["cache_miss"] / PER_MILLION * pricing.chat_cache_miss_per_m
hit_rate = _effective_rate(
pricing.chat_cache_hit_per_m,
pricing.chat_cache_hit_per_m_tier2,
use_tier2,
pricing,
off_peak,
)
output_cost = norm["completion"] / PER_MILLION * pricing.chat_output_per_m
miss_rate = _effective_rate(
pricing.chat_cache_miss_per_m,
pricing.chat_cache_miss_per_m_tier2,
use_tier2,
pricing,
off_peak,
)
output_rate = _effective_rate(
pricing.chat_output_per_m,
pricing.chat_output_per_m_tier2,
use_tier2,
pricing,
off_peak,
)
input_cost = (
norm["cache_hit"] / PER_MILLION * hit_rate
+ norm["cache_miss"] / PER_MILLION * miss_rate
)
output_cost = norm["completion"] / PER_MILLION * output_rate
native = usage.get("cost") if isinstance(usage, dict) else None
if isinstance(native, (int, float)) and not isinstance(native, bool):
total = max(0.0, float(native))
+3
View File
@@ -7,6 +7,7 @@ from typing import Optional
from devplacepy.config import PRESENCE_ONLINE_LIMIT
from devplacepy.database import db, get_users_by_uids
from devplacepy.database.awards import award_is_prominent
from devplacepy.services import presence
from devplacepy.services.base import BaseService
from devplacepy.services.pubsub import publish as pubsub_publish
@@ -104,6 +105,8 @@ class PresenceRelayService(BaseService):
"uid": row["uid"],
"username": row["username"],
"avatar_seed": row.get("avatar_seed") or row["username"],
"last_award_slug": row.get("last_award_slug") or "",
"award_prominent": award_is_prominent(row),
}
for row in users
],
@@ -0,0 +1 @@
# retoor <retoor@molodetz.nl>
+77
View File
@@ -0,0 +1,77 @@
# retoor <retoor@molodetz.nl>
from __future__ import annotations
from devplacepy.services.statistics.common import (
VALID_TABS,
cache_get,
cache_set,
iso,
now_utc,
window_config,
)
from devplacepy.services.statistics.tabs_data import (
build_ai,
build_awards,
build_containers,
build_content,
build_devii,
build_engagement,
build_game,
build_members,
build_moderation,
build_overview,
build_services,
build_social,
build_storage,
build_tools,
build_visitors,
)
TAB_BUILDERS = {
"overview": build_overview,
"visitors": build_visitors,
"members": build_members,
"content": build_content,
"engagement": build_engagement,
"social": build_social,
"ai": build_ai,
"devii": build_devii,
"services": build_services,
"containers": build_containers,
"game": build_game,
"awards": build_awards,
"moderation": build_moderation,
"tools": build_tools,
"storage": build_storage,
}
def build_statistics_tab(
tab: str,
hours: int = 168,
*,
compare: bool = True,
top_n: int = 10,
) -> dict:
key = (tab or "overview").strip().lower()
if key not in VALID_TABS:
key = "overview"
bounded_top = max(1, min(int(top_n or 10), 50))
window = window_config(hours)
cache_key = f"stats:{key}:{window['hours']}:{int(compare)}:{bounded_top}"
cached = cache_get(cache_key)
if cached is not None:
return cached
builder = TAB_BUILDERS[key]
body = builder(window["hours"], compare, bounded_top)
payload = {
"tab": key,
"window_hours": window["hours"],
"granularity": window["granularity"],
"generated_at": iso(now_utc()),
"compare": bool(compare),
**body,
}
cache_set(cache_key, payload)
return payload
+169
View File
@@ -0,0 +1,169 @@
# retoor <retoor@molodetz.nl>
from __future__ import annotations
from datetime import datetime, timedelta, timezone
from typing import Optional
from devplacepy.cache import TTLCache
MAX_WINDOW_HOURS = 2160
RETENTION_DAYS = 90
_stats_cache = TTLCache(ttl=60, max_size=200)
VALID_TABS = (
"overview",
"visitors",
"members",
"content",
"engagement",
"social",
"ai",
"devii",
"services",
"containers",
"game",
"awards",
"moderation",
"tools",
"storage",
)
def now_utc() -> datetime:
return datetime.now(timezone.utc)
def iso(moment: datetime) -> str:
return moment.isoformat(timespec="seconds")
def hour_bucket(moment: datetime) -> str:
return moment.replace(minute=0, second=0, microsecond=0).isoformat(timespec="seconds")
def window_config(hours: int) -> dict:
now = now_utc()
if hours <= 0:
return {
"hours": 0,
"start": None,
"compare_start": None,
"end": iso(now),
"granularity": "weekly",
}
bounded = max(1, min(int(hours), MAX_WINDOW_HOURS))
start = now - timedelta(hours=bounded)
compare_start = start - timedelta(hours=bounded)
if bounded <= 48:
granularity = "hourly"
elif bounded <= 720:
granularity = "daily"
else:
granularity = "daily"
return {
"hours": bounded,
"start": iso(start),
"compare_start": iso(compare_start),
"end": iso(now),
"granularity": granularity,
}
def bucket_expr(granularity: str, column: str = "created_at") -> str:
if granularity == "hourly":
return f"strftime('%Y-%m-%dT%H:00:00', {column})"
if granularity == "weekly":
return f"strftime('%Y-W%W', {column})"
return f"date({column})"
def delta_pair(current: float, previous: float) -> dict:
if previous == 0:
pct = 100.0 if current > 0 else 0.0
else:
pct = round((current - previous) / previous * 100, 1)
if pct > 0:
direction = "up"
elif pct < 0:
direction = "down"
else:
direction = "flat"
return {"delta": pct, "direction": direction}
def metric_card(
key: str,
label: str,
value,
*,
current: Optional[float] = None,
previous: Optional[float] = None,
format_kind: str = "int",
) -> dict:
card = {"key": key, "label": label, "value": value, "format": format_kind}
if current is not None and previous is not None:
card.update(delta_pair(current, previous))
return card
def fill_series(
points: dict[str, float],
start: Optional[str],
end: str,
granularity: str,
) -> list[dict]:
if not start:
ordered = sorted(points.items())
return [{"t": bucket, "v": points[bucket]} for bucket, _ in ordered]
try:
cursor = datetime.fromisoformat(start)
end_dt = datetime.fromisoformat(end)
except (TypeError, ValueError):
ordered = sorted(points.items())
return [{"t": bucket, "v": points[bucket]} for bucket, _ in ordered]
if cursor.tzinfo is None:
cursor = cursor.replace(tzinfo=timezone.utc)
if end_dt.tzinfo is None:
end_dt = end_dt.replace(tzinfo=timezone.utc)
out: list[dict] = []
if granularity == "hourly":
step = timedelta(hours=1)
fmt = lambda moment: hour_bucket(moment)
elif granularity == "weekly":
step = timedelta(days=7)
fmt = lambda moment: moment.strftime("%Y-W%W")
else:
step = timedelta(days=1)
fmt = lambda moment: moment.date().isoformat()
while cursor <= end_dt:
key = fmt(cursor)
out.append({"t": key, "v": points.get(key, 0)})
cursor += step
return out
def cache_get(key: str):
return _stats_cache.get(key)
def cache_set(key: str, payload: dict) -> None:
_stats_cache.set(key, payload)
def series_payload(
key: str,
label: str,
points: dict[str, float],
window: dict,
) -> dict:
return {
"key": key,
"label": label,
"points": fill_series(points, window["start"], window["end"], window["granularity"]),
}
def table_payload(key: str, title: str, columns: list[str], rows: list[list]) -> dict:
return {"key": key, "title": title, "columns": columns, "rows": rows}
File diff suppressed because it is too large Load Diff
+207
View File
@@ -0,0 +1,207 @@
# retoor <retoor@molodetz.nl>
from __future__ import annotations
import asyncio
import hashlib
import logging
from collections import defaultdict
from datetime import datetime, timedelta, timezone
from typing import Optional
from urllib.parse import urlparse
from devplacepy.config import SECRET_KEY
from devplacepy.database import db, get_setting
from devplacepy.utils import client_ip, get_current_user
from .common import RETENTION_DAYS, hour_bucket, now_utc
logger = logging.getLogger(__name__)
EXCLUDED_PREFIXES = (
"/static",
"/avatar",
"/openai",
"/devii/ws",
"/xmlrpc",
"/dbapi",
"/admin/statistics/data",
"/admin/ai-usage/data",
"/admin/services/data",
)
FLUSH_SECONDS = 30
_view_buffer: dict[tuple[str, str, str], dict[str, int]] = defaultdict(
lambda: {"views": 0, "member_views": 0, "guest_views": 0}
)
_unique_buffer: list[tuple[str, str, str, Optional[str]]] = []
_cached_daily_salt = ""
_flush_task: Optional[asyncio.Task] = None
def _tracking_enabled() -> bool:
return get_setting("statistics_tracking_enabled", "1") == "1"
def _daily_salt() -> str:
global _cached_daily_salt
day = now_utc().date().isoformat()
if not _cached_daily_salt or not _cached_daily_salt.startswith(day):
_cached_daily_salt = f"{day}:{SECRET_KEY}"
return _cached_daily_salt
def _visitor_hash(request) -> str:
ip = client_ip(request, default="") or ""
ua = (request.headers.get("user-agent") or "")[:120]
raw = f"{_daily_salt()}:{ip}:{ua}"
return hashlib.sha256(raw.encode("utf-8")).hexdigest()
def _page_group(request) -> str:
route = request.scope.get("route")
path = getattr(route, "path", None)
if path:
return path
segments = [segment for segment in request.url.path.split("/") if segment]
if not segments:
return "/"
return f"/{segments[0]}"
def _referrer_group(request) -> str:
raw = request.headers.get("referer") or request.headers.get("referrer") or ""
if not raw:
return "direct"
try:
host = (urlparse(raw).hostname or "").lower()
except ValueError:
return "other"
if not host:
return "direct"
request_host = (request.url.hostname or "").lower()
if host == request_host:
return "internal"
return host[:120]
def should_track(request, status_code: int) -> bool:
if not _tracking_enabled():
return False
if request.method != "GET":
return False
if status_code >= 400:
return False
path = request.url.path
for prefix in EXCLUDED_PREFIXES:
if path.startswith(prefix):
return False
accept = (request.headers.get("accept") or "").lower()
if "application/json" in accept and "text/html" not in accept:
return False
return True
def track_visit(request, status_code: int) -> None:
if not should_track(request, status_code):
return
bucket = hour_bucket(now_utc())
page_group = _page_group(request)
referrer_group = _referrer_group(request)
key = (bucket, page_group, referrer_group)
slot = _view_buffer[key]
slot["views"] += 1
user = get_current_user(request)
if user:
slot["member_views"] += 1
member_uid = user["uid"]
else:
slot["guest_views"] += 1
member_uid = None
visitor_hash = _visitor_hash(request)
_unique_buffer.append((bucket, visitor_hash, page_group, member_uid))
def _flush_views() -> None:
if "visit_stats_hourly" not in db.tables:
return
if not _view_buffer:
return
pending = dict(_view_buffer)
_view_buffer.clear()
with db:
for (bucket, page_group, referrer_group), counts in pending.items():
db.query(
"INSERT INTO visit_stats_hourly "
"(bucket_start, page_group, referrer_group, views, member_views, guest_views) "
"VALUES (:bucket, :page, :ref, :views, :members, :guests) "
"ON CONFLICT(bucket_start, page_group, referrer_group) DO UPDATE SET "
"views = views + excluded.views, "
"member_views = member_views + excluded.member_views, "
"guest_views = guest_views + excluded.guest_views",
bucket=bucket,
page=page_group,
ref=referrer_group,
views=counts["views"],
members=counts["member_views"],
guests=counts["guest_views"],
)
def _flush_uniques() -> None:
if "visit_unique_slots" not in db.tables:
return
if not _unique_buffer:
return
pending = list(_unique_buffer)
_unique_buffer.clear()
with db:
for bucket, visitor_hash, page_group, member_uid in pending:
db.query(
"INSERT OR IGNORE INTO visit_unique_slots "
"(bucket_start, visitor_hash, page_group, user_uid) "
"VALUES (:bucket, :hash, :page, :uid)",
bucket=bucket,
hash=visitor_hash,
page=page_group,
uid=member_uid,
)
def _prune_old() -> None:
cutoff = (now_utc() - timedelta(days=RETENTION_DAYS)).isoformat(timespec="seconds")
with db:
if "visit_stats_hourly" in db.tables:
db.query("DELETE FROM visit_stats_hourly WHERE bucket_start < :cutoff", cutoff=cutoff)
if "visit_unique_slots" in db.tables:
db.query("DELETE FROM visit_unique_slots WHERE bucket_start < :cutoff", cutoff=cutoff)
def flush_visits(*, prune: bool = False) -> None:
try:
_flush_views()
_flush_uniques()
if prune:
_prune_old()
except Exception:
logger.exception("visit statistics flush failed")
async def _flush_loop() -> None:
ticks = 0
while True:
await asyncio.sleep(FLUSH_SECONDS)
ticks += 1
flush_visits(prune=ticks % 120 == 0)
def start_visit_flusher() -> None:
global _flush_task
if _flush_task is not None:
return
try:
loop = asyncio.get_running_loop()
except RuntimeError:
return
_flush_task = loop.create_task(_flush_loop())