wattgpu / wattgpu_demo /demand_log.py
maufadel's picture
Fixed persistent storage, issue with requirements
05985dc
Raw History Blame Contribute Delete
6.4 kB
"""Record which (model, GPU) pairs the demo cannot answer.
Every refusal is a request someone wanted and did not get, so the log doubles
as a demand-ranked coverage roadmap: it says which models and which GPUs are
worth measuring or supporting next.
Only the query itself is written -- the model id, the GPU name, the scenario
and why it was refused. No identifiers, no addresses, nothing about who asked,
and never a visitor's access token. The UI discloses that this happens.
Persistence
-----------
A Hugging Face Space has an ephemeral filesystem: it is wiped whenever the
Space restarts, which for a quiet Space is often. Set two environment variables
and the log is mirrored to a Hugging Face dataset repository instead, using
`CommitScheduler` from `huggingface_hub` (already a Gradio dependency):
WATTGPU_DEMAND_DATASET=<user>/wattgpu-demand # may be a private repo
WATTGPU_DATASET_TOKEN=hf_... # write access to that repo
Writes still go to a local file; the scheduler batches them up to the dataset
every few minutes, so a restart loses at most one interval. Without those
variables nothing changes and the log stays local.
"""
from __future__ import annotations
import collections
import glob
import json
import os
import threading
import uuid
from datetime import datetime, timezone
DATA_DIR = os.path.join(os.path.dirname(os.path.dirname(os.path.abspath(__file__))), "data")
# The log lives in a folder of its own: the sync uploads whatever is in it, and
# `data/` also holds the models and the GPU database, which must not be pushed.
DEFAULT_LOG_DIR = os.path.join(DATA_DIR, "demand")
# Several Space replicas would otherwise overwrite each other's file, since the
# scheduler uploads the folder wholesale. One shard per process avoids that.
INSTANCE_ID = uuid.uuid4().hex[:8]
# Appends are short, but several viewers can hit the app at once.
_LOCK = threading.Lock()
# Truncated so a pathological input cannot bloat the log.
MAX_FIELD_CHARS = 200
# How often to push to the dataset, in minutes.
DEFAULT_SYNC_MINUTES = 5.0
_scheduler = None
def log_dir() -> str:
return os.environ.get("WATTGPU_DEMAND_LOG_DIR", DEFAULT_LOG_DIR)
def log_path() -> str:
"""The file this process appends to.
`WATTGPU_DEMAND_LOG` still names a file directly, which keeps existing
deployments and the tests working.
"""
explicit = os.environ.get("WATTGPU_DEMAND_LOG")
if explicit:
return explicit
suffix = f"-{INSTANCE_ID}" if is_syncing() else ""
return os.path.join(log_dir(), f"demand_log{suffix}.jsonl")
def is_enabled() -> bool:
"""Logging is on unless explicitly disabled."""
return os.environ.get("WATTGPU_DEMAND_LOG_DISABLED", "").lower() not in ("1", "true", "yes")
def is_syncing() -> bool:
"""Whether a dataset repository is configured to mirror the log to."""
return bool(os.environ.get("WATTGPU_DEMAND_DATASET"))
def record_refusal(model_id: str, gpu_name: str, scenario: str, reason: str) -> None:
"""Append one refused query. Never raises: logging must not break a response."""
if not is_enabled():
return
entry = {
"at": datetime.now(timezone.utc).isoformat(timespec="seconds"),
"model": str(model_id or "")[:MAX_FIELD_CHARS],
"gpu": str(gpu_name or "")[:MAX_FIELD_CHARS],
"scenario": str(scenario or "")[:MAX_FIELD_CHARS],
"reason": reason,
}
try:
path = log_path()
directory = os.path.dirname(path)
if directory:
os.makedirs(directory, exist_ok=True)
with _LOCK, open(path, "a", encoding="utf-8") as fh:
fh.write(json.dumps(entry) + "\n")
except OSError:
pass # a read-only or full filesystem must not take the demo down
def start_sync() -> str | None:
"""Begin mirroring the log folder to a dataset repository.
Returns a short status line for the startup log, or None when no dataset is
configured. Failures are reported and swallowed: the demo must still serve
estimates if the sync cannot start.
"""
global _scheduler
repo_id = os.environ.get("WATTGPU_DEMAND_DATASET")
if not (repo_id and is_enabled()):
return None
if _scheduler is not None:
return f"demand log already syncing to {repo_id}"
token = os.environ.get("WATTGPU_DATASET_TOKEN") or os.environ.get("HF_TOKEN")
if not token:
return ("demand log NOT syncing: WATTGPU_DEMAND_DATASET is set but no "
"WATTGPU_DATASET_TOKEN (or HF_TOKEN) to write with")
try:
minutes = float(os.environ.get("WATTGPU_SYNC_MINUTES", DEFAULT_SYNC_MINUTES))
except ValueError:
minutes = DEFAULT_SYNC_MINUTES
try:
from huggingface_hub import CommitScheduler
os.makedirs(log_dir(), exist_ok=True)
_scheduler = CommitScheduler(
repo_id=repo_id,
repo_type="dataset",
folder_path=log_dir(),
every=minutes,
token=token,
private=True, # only applies when the repo is created here
allow_patterns=["*.jsonl"],
)
return f"demand log syncing to dataset {repo_id} every {minutes:g} min"
except Exception as exc: # noqa: BLE001 - never fail startup over telemetry
return f"demand log NOT syncing ({type(exc).__name__}: {exc})"
def summarise(path: str | None = None, limit: int = 20) -> list[tuple[str, int]]:
"""Most-requested refused queries, for deciding what to cover next.
Reads a single file when given one, otherwise every shard in the log folder.
"""
if path:
sources = [path]
else:
explicit = os.environ.get("WATTGPU_DEMAND_LOG")
sources = [explicit] if explicit else sorted(
glob.glob(os.path.join(log_dir(), "*.jsonl")))
counts: collections.Counter[str] = collections.Counter()
for source in sources:
try:
with open(source, encoding="utf-8") as fh:
for line in fh:
try:
entry = json.loads(line)
except json.JSONDecodeError:
continue
counts[f"{entry.get('reason')}: {entry.get('model')} @ "
f"{entry.get('gpu')}"] += 1
except OSError:
continue
return counts.most_common(limit)