vadugwi / app /persistence.py
deucebucket's picture
Data flywheel: 50q + word ratings + FastAPI backend (anonymous collection)
6da6c49 verified
Raw
History Blame Contribute Delete
3.69 kB
"""Anonymous data collection via append-only JSONL + throttled HF Dataset push.
Modeled on the clanker-pet Snapshotter: buffer rows in memory, append to a local
JSONL on disk, and (throttled, background, fail-soft) push the JSONL to the PRIVATE
dataset deucebucket/vadugwi-data using huggingface_hub + os.environ HF_TOKEN.
Nothing here EVER raises into the caller. If HF_TOKEN is absent (local dev) we just
write the local JSONL and skip the push (logged once). A dataset hiccup must never
break /api/submit.
PRIVACY (HARD, design spec 11.3): rows carry ONLY {ts, engine_version, nonce, ...
rating fields}. No usernames, device ids, cookies, IPs, headers, or raw probe text.
"""
from __future__ import annotations
import json
import os
import threading
import time
STATE_FILE = "vadugwi-data.jsonl"
def _token() -> str | None:
return os.environ.get("HF_TOKEN")
class DataSink:
"""Append-only anonymous row sink: local JSONL + throttled fail-soft HF push."""
def __init__(
self,
jsonl_path: str,
repo_id: str,
min_interval: float = 60.0,
now=time.monotonic,
):
self._path = jsonl_path
self._repo = repo_id
self._min = min_interval
self._now = now
self._last = 0.0
self._lock = threading.Lock()
self._thread: threading.Thread | None = None
self._warned_no_token = False
os.makedirs(os.path.dirname(self._path) or ".", exist_ok=True)
def append(self, rows: list[dict]) -> None:
"""Append rows to the local JSONL (always), then maybe push. Fail-soft."""
if not rows:
return
try:
with open(self._path, "a", encoding="utf-8") as f:
for row in rows:
f.write(json.dumps(row, ensure_ascii=False) + "\n")
except Exception:
return
self._maybe_push()
def _maybe_push(self) -> None:
if not _token():
if not self._warned_no_token:
self._warned_no_token = True
print("[persistence] HF_TOKEN absent — writing local JSONL only, "
"skipping dataset push (local dev mode).")
return
with self._lock:
t = self._now()
if t - self._last < self._min:
return
self._last = t
self._thread = threading.Thread(target=self._push, daemon=True)
self._thread.start()
def _push(self) -> None:
token = _token()
if not token or not os.path.exists(self._path):
return
try:
from huggingface_hub import HfApi
HfApi().upload_file(
path_or_fileobj=self._path,
path_in_repo=STATE_FILE,
repo_id=self._repo,
repo_type="dataset",
token=token,
)
except Exception:
# dataset hiccup must never surface; next append retries (throttled)
return
def restore(self) -> bool:
"""Optional restore-on-boot: pull the existing JSONL so we keep appending
to the full history instead of starting empty. Fail-soft / best-effort."""
token = _token()
if not token:
return False
try:
from huggingface_hub import hf_hub_download
cached = hf_hub_download(
repo_id=self._repo,
repo_type="dataset",
filename=STATE_FILE,
token=token,
)
import shutil
shutil.copyfile(cached, self._path)
return True
except Exception:
return False